mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-08 11:57:51 +00:00
Merge branch 'master' into reopen
This commit is contained in:
commit
97e7f86372
25 changed files with 374 additions and 60 deletions
29
api.go
29
api.go
|
|
@ -655,6 +655,33 @@ func (api *API) Schema(ctx context.Context) []*IndexInfo {
|
|||
return api.holder.limitedSchema()
|
||||
}
|
||||
|
||||
// ApplySchema takes the given schema and applies it across the
|
||||
// cluster (if remote is false), or just to this node (if remote is
|
||||
// true). This is designed for the use case of replicating a schema
|
||||
// from one Pilosa cluster to another which is initially empty. It is
|
||||
// not officially supported in other scenarios and may produce
|
||||
// surprising results.
|
||||
func (api *API) ApplySchema(ctx context.Context, s *Schema, remote bool) error {
|
||||
span, _ := tracing.StartSpanFromContext(ctx, "API.ApplySchema")
|
||||
defer span.Finish()
|
||||
|
||||
if err := api.validate(apiApplySchema); err != nil {
|
||||
return errors.Wrap(err, "validating api method")
|
||||
}
|
||||
|
||||
if !remote {
|
||||
nodes := api.cluster.Nodes()
|
||||
for i, node := range nodes {
|
||||
err := api.server.defaultClient.PostSchema(ctx, &node.URI, s, true)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "forwarding post schema to node %d of %d", i+1, len(nodes))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return api.holder.applySchema(s)
|
||||
}
|
||||
|
||||
// Views returns the views in the given field.
|
||||
func (api *API) Views(ctx context.Context, indexName string, fieldName string) ([]*view, error) {
|
||||
span, _ := tracing.StartSpanFromContext(ctx, "API.Views")
|
||||
|
|
@ -1273,6 +1300,7 @@ const (
|
|||
//apiStatsWithTags // not implemented
|
||||
//apiVersion // not implemented
|
||||
apiViews
|
||||
apiApplySchema
|
||||
)
|
||||
|
||||
var methodsCommon = map[apiMethod]struct{}{
|
||||
|
|
@ -1306,4 +1334,5 @@ var methodsNormal = map[apiMethod]struct{}{
|
|||
apiRemoveNode: {},
|
||||
apiShardNodes: {},
|
||||
apiViews: {},
|
||||
apiApplySchema: {},
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4,9 +4,40 @@ package pilosa
|
|||
|
||||
import "strconv"
|
||||
|
||||
const _apiMethod_name = "apiClusterMessageapiCreateFieldapiCreateIndexapiDeleteFieldapiDeleteAvailableShardapiDeleteIndexapiDeleteViewapiExportCSVapiFragmentBlockDataapiFragmentBlocksapiFieldapiFieldAttrDiffapiImportapiImportValueapiIndexapiIndexAttrDiffapiQueryapiRecalculateCachesapiRemoveNodeapiResizeAbortapiSetCoordinatorapiShardNodesapiViews"
|
||||
func _() {
|
||||
// An "invalid array index" compiler error signifies that the constant values have changed.
|
||||
// Re-run the stringer command to generate them again.
|
||||
var x [1]struct{}
|
||||
_ = x[apiClusterMessage-0]
|
||||
_ = x[apiCreateField-1]
|
||||
_ = x[apiCreateIndex-2]
|
||||
_ = x[apiDeleteField-3]
|
||||
_ = x[apiDeleteAvailableShard-4]
|
||||
_ = x[apiDeleteIndex-5]
|
||||
_ = x[apiDeleteView-6]
|
||||
_ = x[apiExportCSV-7]
|
||||
_ = x[apiFragmentBlockData-8]
|
||||
_ = x[apiFragmentBlocks-9]
|
||||
_ = x[apiFragmentData-10]
|
||||
_ = x[apiField-11]
|
||||
_ = x[apiFieldAttrDiff-12]
|
||||
_ = x[apiImport-13]
|
||||
_ = x[apiImportValue-14]
|
||||
_ = x[apiIndex-15]
|
||||
_ = x[apiIndexAttrDiff-16]
|
||||
_ = x[apiQuery-17]
|
||||
_ = x[apiRecalculateCaches-18]
|
||||
_ = x[apiRemoveNode-19]
|
||||
_ = x[apiResizeAbort-20]
|
||||
_ = x[apiSetCoordinator-21]
|
||||
_ = x[apiShardNodes-22]
|
||||
_ = x[apiViews-23]
|
||||
_ = x[apiApplySchema-24]
|
||||
}
|
||||
|
||||
var _apiMethod_index = [...]uint16{0, 17, 31, 45, 59, 82, 96, 109, 121, 141, 158, 166, 182, 191, 205, 213, 229, 237, 257, 270, 284, 301, 314, 322}
|
||||
const _apiMethod_name = "apiClusterMessageapiCreateFieldapiCreateIndexapiDeleteFieldapiDeleteAvailableShardapiDeleteIndexapiDeleteViewapiExportCSVapiFragmentBlockDataapiFragmentBlocksapiFragmentDataapiFieldapiFieldAttrDiffapiImportapiImportValueapiIndexapiIndexAttrDiffapiQueryapiRecalculateCachesapiRemoveNodeapiResizeAbortapiSetCoordinatorapiShardNodesapiViewsapiApplySchema"
|
||||
|
||||
var _apiMethod_index = [...]uint16{0, 17, 31, 45, 59, 82, 96, 109, 121, 141, 158, 173, 181, 197, 206, 220, 228, 244, 252, 272, 285, 299, 316, 329, 337, 351}
|
||||
|
||||
func (i apiMethod) String() string {
|
||||
if i < 0 || i >= apiMethod(len(_apiMethod_index)-1) {
|
||||
|
|
|
|||
|
|
@ -46,6 +46,7 @@ type FieldValue struct {
|
|||
type InternalClient interface {
|
||||
MaxShardByIndex(ctx context.Context) (map[string]uint64, error)
|
||||
Schema(ctx context.Context) ([]*IndexInfo, error)
|
||||
PostSchema(ctx context.Context, uri *URI, s *Schema, remote bool) error
|
||||
CreateIndex(ctx context.Context, index string, opt IndexOptions) error
|
||||
FragmentNodes(ctx context.Context, index string, shard uint64) ([]*Node, error)
|
||||
Nodes(ctx context.Context) ([]*Node, error)
|
||||
|
|
@ -103,6 +104,10 @@ func (n nopInternalClient) MaxShardByIndex(context.Context) (map[string]uint64,
|
|||
return nil, nil
|
||||
}
|
||||
func (n nopInternalClient) Schema(ctx context.Context) ([]*IndexInfo, error) { return nil, nil }
|
||||
func (n nopInternalClient) PostSchema(ctx context.Context, uri *URI, s *Schema, remote bool) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (n nopInternalClient) CreateIndex(ctx context.Context, index string, opt IndexOptions) error {
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -44,7 +44,9 @@ func TestExportCommand_Validation(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestExportCommand_Run(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
|
||||
buf := bytes.Buffer{}
|
||||
stdin, stdout, stderr := GetIO(buf)
|
||||
|
|
|
|||
|
|
@ -70,7 +70,9 @@ func TestImportCommand_Basic(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
cm.Host = cmd.API.Node().URI.HostPort()
|
||||
|
||||
cm.Index = "i"
|
||||
|
|
@ -97,7 +99,9 @@ func TestImportCommand_Basic(t *testing.T) {
|
|||
}
|
||||
ctx := context.Background()
|
||||
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
cm.Host = cmd.API.Node().URI.HostPort()
|
||||
|
||||
cm.Index = "i"
|
||||
|
|
@ -128,7 +132,9 @@ func TestImportCommand_RunValue(t *testing.T) {
|
|||
}
|
||||
ctx := context.Background()
|
||||
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
cm.Host = cmd.API.Node().URI.HostPort()
|
||||
|
||||
resp, err := http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i", strings.NewReader("")))
|
||||
|
|
@ -168,7 +174,9 @@ func TestImportCommand_RunValue(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
cm.Host = cmd.API.Node().URI.HostPort()
|
||||
|
||||
resp, err := http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i", strings.NewReader("")))
|
||||
|
|
@ -208,7 +216,9 @@ func TestImportCommand_RunKeys(t *testing.T) {
|
|||
}
|
||||
ctx := context.Background()
|
||||
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
cm.Host = cmd.API.Node().URI.HostPort()
|
||||
|
||||
resp, err := http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i", strings.NewReader(`{"options":{"keys": true}}`)))
|
||||
|
|
@ -259,6 +269,7 @@ func TestImportCommand_KeyReplication(t *testing.T) {
|
|||
ctx := context.Background()
|
||||
|
||||
c := test.MustRunCluster(t, 2)
|
||||
defer c.Close()
|
||||
cmd0 := c[0]
|
||||
cmd1 := c[1]
|
||||
|
||||
|
|
@ -319,7 +330,9 @@ func TestImportCommand_RunValueKeys(t *testing.T) {
|
|||
}
|
||||
ctx := context.Background()
|
||||
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
cm.Host = cmd.API.Node().URI.HostPort()
|
||||
|
||||
resp, err := http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i", strings.NewReader(`{"options":{"keys": true}}`)))
|
||||
|
|
@ -343,7 +356,9 @@ func TestImportCommand_RunValueKeys(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestImportCommand_InvalidFile(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
|
||||
buf := bytes.Buffer{}
|
||||
stdin, stdout, stderr := GetIO(buf)
|
||||
|
|
@ -429,7 +444,9 @@ func GetIO(buf bytes.Buffer) (io.Reader, io.Writer, io.Writer) {
|
|||
}
|
||||
|
||||
func TestImportCommand_BugOverwriteValue(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
|
||||
buf := bytes.Buffer{}
|
||||
stdin, stdout, stderr := GetIO(buf)
|
||||
|
|
@ -503,7 +520,9 @@ func TestImportCommand_RunBool(t *testing.T) {
|
|||
cm := NewImportCommand(stdin, stdout, stderr)
|
||||
ctx := context.Background()
|
||||
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
cm.Host = cmd.API.Node().URI.HostPort()
|
||||
|
||||
resp, err := http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i", strings.NewReader("")))
|
||||
|
|
|
|||
|
|
@ -241,7 +241,7 @@ curl -XDELETE localhost:10101/index/user/field/language
|
|||
Returns the schema of all indexes in JSON.
|
||||
|
||||
``` request
|
||||
curl -XGET localhost:10101/index
|
||||
curl -XGET localhost:10101/schema
|
||||
```
|
||||
``` response
|
||||
{
|
||||
|
|
@ -276,6 +276,24 @@ curl -XGET localhost:10101/index
|
|||
}
|
||||
```
|
||||
|
||||
### Duplicate schema into empty Pilosa cluster
|
||||
|
||||
`POST /schema`
|
||||
|
||||
To duplicate one Pilosa cluster's schema to another, it's possible to
|
||||
pass the output of `GET /schema` as the request body of `POST /schema`
|
||||
and all the indexes and fields in the schema will be created in
|
||||
Pilosa. As of this writing, the behavior of POSTing a schema to a
|
||||
non-empty Pilosa cluster is undefined. These semantics will likely be
|
||||
ironed out in a future version.
|
||||
|
||||
``` request
|
||||
# after (e.g.) curl -XGET localhost:10101/schema > schema.json
|
||||
curl -XPOST localhost:10101/schema --data-binary @schema.json
|
||||
```
|
||||
|
||||
Response: `204 No Content`
|
||||
|
||||
### Get version
|
||||
|
||||
`GET /version`
|
||||
|
|
|
|||
|
|
@ -186,7 +186,7 @@ Set(3, A=8, 2017-05-19T00:00)
|
|||
```
|
||||
|
||||

|
||||
*Time quantum fueld diagram*
|
||||
*Time quantum field diagram*
|
||||
|
||||
#### Mutex
|
||||
|
||||
|
|
|
|||
|
|
@ -74,7 +74,7 @@ curl localhost:10101/index/repository/field/stargazer \
|
|||
{"success":true}
|
||||
```
|
||||
|
||||
Since our data contains time stamps whcih represent the time users starred repos, we set the field type to `time`. Time quantum is the resolution of the time we want to use, and we set it to `YMD` (year, month, day) for `stargazer`.
|
||||
Since our data contains time stamps which represent the time users starred repos, we set the field type to `time`. Time quantum is the resolution of the time we want to use, and we set it to `YMD` (year, month, day) for `stargazer`.
|
||||
|
||||
Next up is the `language` field, which will contain IDs for programming languages:
|
||||
``` request
|
||||
|
|
|
|||
|
|
@ -508,7 +508,9 @@ func TestExecutor_Execute_Count(t *testing.T) {
|
|||
// Ensure a set query can be executed.
|
||||
func TestExecutor_Execute_Set(t *testing.T) {
|
||||
t.Run("RowIDColumnID", func(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
holder := cmd.Server.Holder()
|
||||
hldr := test.Holder{Holder: holder}
|
||||
hldr.SetBit("i", "f", 1, 0)
|
||||
|
|
@ -571,7 +573,9 @@ func TestExecutor_Execute_Set(t *testing.T) {
|
|||
})
|
||||
|
||||
t.Run("RowKeyColumnKey", func(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
holder := cmd.Server.Holder()
|
||||
hldr := test.Holder{Holder: holder}
|
||||
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{Keys: true})
|
||||
|
|
|
|||
9
field.go
9
field.go
|
|
@ -379,27 +379,33 @@ func (f *Field) Options() FieldOptions {
|
|||
func (f *Field) Open() error {
|
||||
if err := func() error {
|
||||
// Ensure the field's path exists.
|
||||
f.logger.Debugf("ensure field path exists: %s", f.path)
|
||||
if err := os.MkdirAll(f.path, 0777); err != nil {
|
||||
return errors.Wrap(err, "creating field dir")
|
||||
}
|
||||
|
||||
f.logger.Debugf("load meta file for index/field: %s/%s", f.index, f.name)
|
||||
if err := f.loadMeta(); err != nil {
|
||||
return errors.Wrap(err, "loading meta")
|
||||
}
|
||||
|
||||
f.logger.Debugf("load available shards for index/field: %s/%s", f.index, f.name)
|
||||
if err := f.loadAvailableShards(); err != nil {
|
||||
return errors.Wrap(err, "loading available shards")
|
||||
}
|
||||
|
||||
// Apply the field options loaded from meta.
|
||||
f.logger.Debugf("apply options for index/field: %s/%s", f.index, f.name)
|
||||
if err := f.applyOptions(f.options); err != nil {
|
||||
return errors.Wrap(err, "applying options")
|
||||
}
|
||||
|
||||
f.logger.Debugf("open views for index/field: %s/%s", f.index, f.name)
|
||||
if err := f.openViews(); err != nil {
|
||||
return errors.Wrap(err, "opening views")
|
||||
}
|
||||
|
||||
f.logger.Debugf("open row attribute store for index/field: %s/%s", f.index, f.name)
|
||||
if err := f.rowAttrStore.Open(); err != nil {
|
||||
return errors.Wrap(err, "opening attrstore")
|
||||
}
|
||||
|
|
@ -410,6 +416,7 @@ func (f *Field) Open() error {
|
|||
return err
|
||||
}
|
||||
|
||||
f.logger.Debugf("successfully opened field index/field: %s/%s", f.index, f.name)
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -434,11 +441,13 @@ func (f *Field) openViews() error {
|
|||
}
|
||||
|
||||
name := filepath.Base(fi.Name())
|
||||
f.logger.Debugf("open index/field/view: %s/%s/%s", f.index, f.name, fi.Name())
|
||||
view := f.newView(f.viewPath(name), name)
|
||||
if err := view.open(); err != nil {
|
||||
return fmt.Errorf("opening view: view=%s, err=%s", view.name, err)
|
||||
}
|
||||
view.rowAttrStore = f.rowAttrStore
|
||||
f.logger.Debugf("add index/field/view to field.viewMap: %s/%s/%s", f.index, f.name, view.name)
|
||||
f.viewMap[view.name] = view
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -163,11 +163,13 @@ func (f *fragment) Open() error {
|
|||
|
||||
if err := func() error {
|
||||
// Initialize storage in a function so we can close if anything goes wrong.
|
||||
f.Logger.Debugf("open storage for index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
|
||||
if err := f.openStorage(); err != nil {
|
||||
return errors.Wrap(err, "opening storage")
|
||||
}
|
||||
|
||||
// Fill cache with rows persisted to disk.
|
||||
f.Logger.Debugf("open cache for index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
|
||||
if err := f.openCache(); err != nil {
|
||||
return errors.Wrap(err, "opening cache")
|
||||
}
|
||||
|
|
@ -185,6 +187,7 @@ func (f *fragment) Open() error {
|
|||
return err
|
||||
}
|
||||
|
||||
f.Logger.Debugf("successfully opened index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
|
|||
10
go.mod
10
go.mod
|
|
@ -31,9 +31,13 @@ require (
|
|||
github.com/spf13/viper v1.3.1
|
||||
github.com/uber-go/atomic v1.3.2 // indirect
|
||||
github.com/uber/jaeger-client-go v2.15.0+incompatible
|
||||
github.com/uber/jaeger-lib v1.5.0 // indirect
|
||||
go.uber.org/atomic v1.3.2 // indirect
|
||||
golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4
|
||||
github.com/uber/jaeger-lib v1.5.0
|
||||
golang.org/x/crypto v0.0.0-20190426145343-a29dc8fdc734 // indirect
|
||||
golang.org/x/net v0.0.0-20190424112056-4829fb13d2c6 // indirect
|
||||
golang.org/x/sync v0.0.0-20190423024810-112230192c58
|
||||
golang.org/x/sys v0.0.0-20190429190828-d89cdac9e872 // indirect
|
||||
golang.org/x/text v0.3.2 // indirect
|
||||
golang.org/x/tools v0.0.0-20190429231329-9d4d845e86f1 // indirect
|
||||
modernc.org/mathutil v1.0.0
|
||||
modernc.org/strutil v1.0.0
|
||||
)
|
||||
|
|
|
|||
19
go.sum
19
go.sum
|
|
@ -113,16 +113,33 @@ go.uber.org/atomic v1.3.2/go.mod h1:gD2HeocX3+yG+ygLZcrzQJaqmWj9AIm7n08wl/qW/PE=
|
|||
golang.org/x/crypto v0.0.0-20181029021203-45a5f77698d3/go.mod h1:6SG95UA2DQfeDnfUPMdvaQW0Q7yPrPDi9nlGo2tz2b4=
|
||||
golang.org/x/crypto v0.0.0-20181203042331-505ab145d0a9 h1:mKdxBk7AujPs8kU4m80U72y/zjbZ3UcXC7dClwKbUI0=
|
||||
golang.org/x/crypto v0.0.0-20181203042331-505ab145d0a9/go.mod h1:6SG95UA2DQfeDnfUPMdvaQW0Q7yPrPDi9nlGo2tz2b4=
|
||||
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
|
||||
golang.org/x/crypto v0.0.0-20190426145343-a29dc8fdc734 h1:p/H982KKEjUnLJkM3tt/LemDnOc1GiZL5FCVlORJ5zo=
|
||||
golang.org/x/crypto v0.0.0-20190426145343-a29dc8fdc734/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
|
||||
golang.org/x/net v0.0.0-20181023162649-9b4f9f5ad519 h1:x6rhz8Y9CjbgQkccRGmELH6K+LJj7tOoh3XWeC1yaQM=
|
||||
golang.org/x/net v0.0.0-20181023162649-9b4f9f5ad519/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
|
||||
golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
|
||||
golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
|
||||
golang.org/x/net v0.0.0-20190424112056-4829fb13d2c6 h1:FP8hkuE6yUEaJnK7O2eTuejKWwW+Rhfj80dQ2JcKxCU=
|
||||
golang.org/x/net v0.0.0-20190424112056-4829fb13d2c6/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
|
||||
golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4 h1:YUO/7uOKsKeq9UokNS62b8FYywz3ker1l1vDZRCRefw=
|
||||
golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sync v0.0.0-20190423024810-112230192c58 h1:8gQV6CLnAEikrhgkHFbMAEhagSSnXWGV915qUMm9mrU=
|
||||
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
|
||||
golang.org/x/sys v0.0.0-20181026203630-95b1ffbd15a5/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
|
||||
golang.org/x/sys v0.0.0-20181205085412-a5c9d58dba9a h1:1n5lsVfiQW3yfsRGu98756EH1YthsFqr/5mxHduZW2A=
|
||||
golang.org/x/sys v0.0.0-20181205085412-a5c9d58dba9a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
|
||||
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
|
||||
golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20190429190828-d89cdac9e872 h1:cGjJzUd8RgBw428LXP65YXni0aiGNA4Bl+ls8SmLOm8=
|
||||
golang.org/x/sys v0.0.0-20190429190828-d89cdac9e872/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/text v0.3.0 h1:g61tztE5qeGQ89tm6NTjjM9VPIm088od1l6aSorWRWg=
|
||||
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM=
|
||||
golang.org/x/text v0.3.2 h1:tW2bmiBqwgJj/UpqtC8EpXEZVYOwU0yG4iWbprSVAcs=
|
||||
golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk=
|
||||
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
|
||||
golang.org/x/tools v0.0.0-20190429231329-9d4d845e86f1 h1:MSSXVSCgrxTAYytvleklMKlLdxjexiJWNffJciO1nCI=
|
||||
golang.org/x/tools v0.0.0-20190429231329-9d4d845e86f1/go.mod h1:RgjU9mgBXZiqYHBnxXauZ1Gv1EHHAz9KjViQ78xBX0Q=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
gopkg.in/yaml.v2 v2.2.2 h1:ZCJp+EgiOT7lHqUV2J862kp8Qj64Jo6az82+3Td9dZw=
|
||||
gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
|
||||
|
|
|
|||
|
|
@ -132,6 +132,33 @@ func (c *InternalClient) Schema(ctx context.Context) ([]*pilosa.IndexInfo, error
|
|||
return rsp.Indexes, nil
|
||||
}
|
||||
|
||||
func (c *InternalClient) PostSchema(ctx context.Context, uri *pilosa.URI, s *pilosa.Schema, remote bool) error {
|
||||
u := uri.Path(fmt.Sprintf("/schema?remote=%v", remote))
|
||||
buf, err := json.Marshal(s)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "marshalling schema")
|
||||
}
|
||||
req, err := http.NewRequest("POST", u, bytes.NewReader(buf))
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "creating request")
|
||||
}
|
||||
|
||||
req.Header.Set("Content-Length", strconv.Itoa(len(buf)))
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("Accept", "application/json")
|
||||
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
|
||||
|
||||
resp, err := c.executeRequest(req.WithContext(ctx))
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "executing request")
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode != http.StatusNoContent {
|
||||
return errors.Errorf("unexpected status code: %s", resp.Status)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// CreateIndex creates a new index on the server.
|
||||
func (c *InternalClient) CreateIndex(ctx context.Context, index string, opt pilosa.IndexOptions) error {
|
||||
span, ctx := tracing.StartSpanFromContext(ctx, "InternalClient.CreateIndex")
|
||||
|
|
|
|||
|
|
@ -182,7 +182,10 @@ func TestClient_MultiNode(t *testing.T) {
|
|||
|
||||
// Ensure client can export data.
|
||||
func TestClient_Export(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
|
||||
host := cmd.URL()
|
||||
|
||||
cmd.MustCreateIndex(t, "keyed", pilosa.IndexOptions{Keys: true})
|
||||
|
|
@ -345,7 +348,9 @@ func TestClient_Export(t *testing.T) {
|
|||
|
||||
// Ensure client can bulk import data.
|
||||
func TestClient_Import(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
host := cmd.URL()
|
||||
holder := cmd.Server.Holder()
|
||||
hldr := test.Holder{Holder: holder}
|
||||
|
|
@ -514,7 +519,9 @@ func TestClient_ImportRoaring(t *testing.T) {
|
|||
// Ensure client can bulk import data.
|
||||
func TestClient_ImportKeys(t *testing.T) {
|
||||
t.Run("SingleNode", func(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
host := cmd.URL()
|
||||
|
||||
cmd.MustCreateIndex(t, "keyed", pilosa.IndexOptions{Keys: true})
|
||||
|
|
@ -611,6 +618,7 @@ func TestClient_ImportKeys(t *testing.T) {
|
|||
|
||||
t.Run("MultiNode", func(t *testing.T) {
|
||||
cluster := test.MustRunCluster(t, 2)
|
||||
defer cluster.Close()
|
||||
cmd0 := cluster[0]
|
||||
cmd1 := cluster[1]
|
||||
host0 := cmd0.URL()
|
||||
|
|
@ -686,7 +694,9 @@ func TestClient_ImportKeys(t *testing.T) {
|
|||
})
|
||||
|
||||
t.Run("IntegerFieldSingleNode", func(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
host := cmd.URL()
|
||||
holder := cmd.Server.Holder()
|
||||
hldr := test.Holder{Holder: holder}
|
||||
|
|
@ -769,7 +779,9 @@ func TestClient_ImportKeys(t *testing.T) {
|
|||
|
||||
// Ensure client can bulk import value data.
|
||||
func TestClient_ImportValue(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
host := cmd.URL()
|
||||
holder := cmd.Server.Holder()
|
||||
hldr := test.Holder{Holder: holder}
|
||||
|
|
@ -875,7 +887,9 @@ func TestClient_ImportValue(t *testing.T) {
|
|||
|
||||
// Ensure client can bulk import data while tracking existence.
|
||||
func TestClient_ImportExistence(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
host := cmd.URL()
|
||||
holder := cmd.Server.Holder()
|
||||
hldr := test.Holder{Holder: holder}
|
||||
|
|
@ -952,7 +966,10 @@ func TestClient_ImportExistence(t *testing.T) {
|
|||
|
||||
// Ensure client can retrieve a list of all checksums for blocks in a fragment.
|
||||
func TestClient_FragmentBlocks(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
|
||||
holder := cmd.Server.Holder()
|
||||
hldr := test.Holder{Holder: holder}
|
||||
|
||||
|
|
|
|||
|
|
@ -187,6 +187,7 @@ func (h *Handler) populateValidators() {
|
|||
h.validators["GetInfo"] = queryValidationSpecRequired()
|
||||
h.validators["RecalculateCaches"] = queryValidationSpecRequired()
|
||||
h.validators["GetSchema"] = queryValidationSpecRequired()
|
||||
h.validators["PostSchema"] = queryValidationSpecRequired().Optional("remote")
|
||||
h.validators["GetStatus"] = queryValidationSpecRequired()
|
||||
h.validators["GetVersion"] = queryValidationSpecRequired()
|
||||
h.validators["PostClusterMessage"] = queryValidationSpecRequired()
|
||||
|
|
@ -255,6 +256,7 @@ func newRouter(handler *Handler) *mux.Router {
|
|||
router.HandleFunc("/info", handler.handleGetInfo).Methods("GET").Name("GetInfo")
|
||||
router.HandleFunc("/recalculate-caches", handler.handleRecalculateCaches).Methods("POST").Name("RecalculateCaches")
|
||||
router.HandleFunc("/schema", handler.handleGetSchema).Methods("GET").Name("GetSchema")
|
||||
router.HandleFunc("/schema", handler.handlePostSchema).Methods("POST").Name("PostSchema")
|
||||
router.HandleFunc("/status", handler.handleGetStatus).Methods("GET").Name("GetStatus")
|
||||
router.HandleFunc("/version", handler.handleGetVersion).Methods("GET").Name("GetVersion")
|
||||
|
||||
|
|
@ -412,11 +414,32 @@ func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) {
|
|||
}
|
||||
|
||||
schema := h.api.Schema(r.Context())
|
||||
if err := json.NewEncoder(w).Encode(map[string]interface{}{"indexes": schema}); err != nil {
|
||||
if err := json.NewEncoder(w).Encode(map[string]interface{}{"indexes": schema}); err != nil { // TODO: use pilosa.Schema instead of map[string]interface{} here?
|
||||
h.logger.Printf("write schema response error: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Handler) handlePostSchema(w http.ResponseWriter, r *http.Request) {
|
||||
q := r.URL.Query()
|
||||
remoteStr := q.Get("remote")
|
||||
var remote bool
|
||||
if remoteStr == "true" {
|
||||
remote = true
|
||||
}
|
||||
|
||||
schema := &pilosa.Schema{}
|
||||
if err := json.NewDecoder(r.Body).Decode(schema); err != nil {
|
||||
http.Error(w, fmt.Sprintf("decoding request as JSON Pilosa schema: %v", err), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
if err := h.api.ApplySchema(r.Context(), schema, remote); err != nil {
|
||||
http.Error(w, fmt.Sprintf("apply schema to Pilosa: %v", err), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
}
|
||||
|
||||
// handleGetStatus handles GET /status requests.
|
||||
func (h *Handler) handleGetStatus(w http.ResponseWriter, r *http.Request) {
|
||||
if !validHeaderAcceptJSON(r.Header) {
|
||||
|
|
|
|||
|
|
@ -37,7 +37,9 @@ func TestTranslateStore_Reader(t *testing.T) {
|
|||
// "translator_test.go:65: unexpected EOF"
|
||||
t.Skip()
|
||||
|
||||
primary := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
primary := cluster[0]
|
||||
|
||||
hldr := test.Holder{Holder: primary.Server.Holder()}
|
||||
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{Keys: true})
|
||||
|
|
@ -103,9 +105,10 @@ func TestTranslateStore_Reader(t *testing.T) {
|
|||
}
|
||||
|
||||
opts := server.OptCommandServerOptions(pilosa.OptServerPrimaryTranslateStore(translateStore))
|
||||
primary := test.MustRunCluster(t, 1, []server.CommandOption{opts})[0]
|
||||
cluster := test.MustRunCluster(t, 1, []server.CommandOption{opts})
|
||||
defer cluster.Close()
|
||||
primary := cluster[0]
|
||||
|
||||
defer primary.Close()
|
||||
defer close(done)
|
||||
|
||||
// Connect to server and begin streaming.
|
||||
|
|
@ -135,8 +138,9 @@ func TestTranslateStore_Reader(t *testing.T) {
|
|||
}
|
||||
|
||||
opts := server.OptCommandServerOptions(pilosa.OptServerPrimaryTranslateStore(translateStore))
|
||||
primary := test.MustRunCluster(t, 1, []server.CommandOption{opts})[0]
|
||||
defer primary.Close()
|
||||
cluster := test.MustRunCluster(t, 1, []server.CommandOption{opts})
|
||||
defer cluster.Close()
|
||||
primary := cluster[0]
|
||||
|
||||
ts := http.NewTranslateStore(primary.URL())
|
||||
_, err := ts.Reader(context.Background(), 0)
|
||||
|
|
|
|||
5
index.go
5
index.go
|
|
@ -107,15 +107,18 @@ func (i *Index) options() IndexOptions {
|
|||
// Open opens and initializes the index.
|
||||
func (i *Index) Open() error {
|
||||
// Ensure the path exists.
|
||||
i.logger.Debugf("ensure index path exists: %s", i.path)
|
||||
if err := os.MkdirAll(i.path, 0777); err != nil {
|
||||
return errors.Wrap(err, "creating directory")
|
||||
}
|
||||
|
||||
// Read meta file.
|
||||
i.logger.Debugf("load meta file for index: %s", i.name)
|
||||
if err := i.loadMeta(); err != nil {
|
||||
return errors.Wrap(err, "loading meta file")
|
||||
}
|
||||
|
||||
i.logger.Debugf("open fields for index: %s", i.name)
|
||||
if err := i.openFields(); err != nil {
|
||||
return errors.Wrap(err, "opening fields")
|
||||
}
|
||||
|
|
@ -151,6 +154,7 @@ func (i *Index) openFields() error {
|
|||
continue
|
||||
}
|
||||
|
||||
i.logger.Debugf("open field: %s", fi.Name())
|
||||
fld, err := i.newField(i.fieldPath(filepath.Base(fi.Name())), filepath.Base(fi.Name()))
|
||||
if err != nil {
|
||||
return errors.Wrapf(ErrName, "'%s'", fi.Name())
|
||||
|
|
@ -158,6 +162,7 @@ func (i *Index) openFields() error {
|
|||
if err := fld.Open(); err != nil {
|
||||
return fmt.Errorf("open field: name=%s, err=%s", fld.Name(), err)
|
||||
}
|
||||
i.logger.Debugf("add field to index.fields: %s", fi.Name())
|
||||
i.fields[fld.Name()] = fld
|
||||
}
|
||||
return nil
|
||||
|
|
|
|||
|
|
@ -33,8 +33,7 @@ import (
|
|||
func TestMain_SendReceiveMessage(t *testing.T) {
|
||||
ms := test.MustRunCluster(t, 2)
|
||||
m0, m1 := ms[0], ms[1]
|
||||
defer m0.Close()
|
||||
defer m1.Close()
|
||||
defer ms.Close()
|
||||
|
||||
// Expected indexes and Fields
|
||||
expected := map[string][]string{
|
||||
|
|
@ -127,8 +126,7 @@ func TestClusterResize_EmptyNode(t *testing.T) {
|
|||
// Ensure that a cluster of empty nodes comes up in a NORMAL state.
|
||||
func TestClusterResize_EmptyNodes(t *testing.T) {
|
||||
clus := test.MustRunCluster(t, 2)
|
||||
defer clus[0].Close()
|
||||
defer clus[1].Close()
|
||||
defer clus.Close()
|
||||
|
||||
if clus[0].API.State() != pilosa.ClusterStateNormal {
|
||||
t.Fatalf("unexpected node0 cluster state: %s", clus[0].API.State())
|
||||
|
|
@ -141,6 +139,7 @@ func TestClusterResize_EmptyNodes(t *testing.T) {
|
|||
func TestClusterResize_AddNode(t *testing.T) {
|
||||
t.Run("NoData", func(t *testing.T) {
|
||||
clus := test.MustRunCluster(t, 2)
|
||||
defer clus.Close()
|
||||
|
||||
if !checkClusterState(clus[0], pilosa.ClusterStateNormal, 1000) {
|
||||
t.Fatalf("unexpected node0 cluster state: %s", clus[0].API.State())
|
||||
|
|
@ -552,6 +551,7 @@ func TestCluster_GossipMembership(t *testing.T) {
|
|||
|
||||
func TestClusterResize_RemoveNode(t *testing.T) {
|
||||
cluster := test.MustRunCluster(t, 3)
|
||||
defer cluster.Close()
|
||||
m0 := cluster[0]
|
||||
m1 := cluster[1]
|
||||
|
||||
|
|
|
|||
|
|
@ -16,6 +16,7 @@ package server
|
|||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net"
|
||||
"os"
|
||||
"strings"
|
||||
|
|
@ -125,29 +126,31 @@ func TestConfig_validateAddrs(t *testing.T) {
|
|||
}
|
||||
|
||||
for i, test := range tests {
|
||||
c := NewConfig()
|
||||
t.Run(fmt.Sprintf("%d", i), func(t *testing.T) {
|
||||
c := NewConfig()
|
||||
|
||||
c.Bind = test.in.bind
|
||||
c.Advertise = test.in.advertise
|
||||
c.Bind = test.in.bind
|
||||
c.Advertise = test.in.advertise
|
||||
|
||||
err := c.validateAddrs(context.Background())
|
||||
err := c.validateAddrs(context.Background())
|
||||
|
||||
if err != nil && test.expErr == "" {
|
||||
t.Fatal(err)
|
||||
} else if err == nil && test.expErr != "" {
|
||||
t.Fatalf("test %d: expected error string to contain %s, but got no error", i, test.expErr)
|
||||
} else if err != nil && test.expErr != "" {
|
||||
if strings.Contains(err.Error(), test.expErr) {
|
||||
continue
|
||||
} else {
|
||||
t.Fatalf("test %d: expected error string to contain %s, but got %s", i, test.expErr, err.Error())
|
||||
if err != nil && test.expErr == "" {
|
||||
t.Fatal(err)
|
||||
} else if err == nil && test.expErr != "" {
|
||||
t.Fatalf("expected error string to contain %s, but got no error", test.expErr)
|
||||
} else if err != nil && test.expErr != "" {
|
||||
if strings.Contains(err.Error(), test.expErr) {
|
||||
return
|
||||
} else {
|
||||
t.Fatalf("expected error string to contain %s, but got %s", test.expErr, err.Error())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if c.Bind != test.exp.bind {
|
||||
t.Fatalf("test %d: bind address: expected %s, but got %s", i, test.exp.bind, c.Bind)
|
||||
} else if c.Advertise != test.exp.advertise {
|
||||
t.Fatalf("test %d: advertise address: expected %s, but got %s", i, test.exp.advertise, c.Advertise)
|
||||
}
|
||||
if c.Bind != test.exp.bind {
|
||||
t.Fatalf("bind address: expected %s, but got %s", test.exp.bind, c.Bind)
|
||||
} else if c.Advertise != test.exp.advertise {
|
||||
t.Fatalf("advertise address: expected %s, but got %s", test.exp.advertise, c.Advertise)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -36,13 +36,54 @@ import (
|
|||
"github.com/pilosa/pilosa/test"
|
||||
)
|
||||
|
||||
// Ensure the handler returns "not found" for invalid paths.
|
||||
func TestHandler_PostSchemaCluster(t *testing.T) {
|
||||
cluster := test.MustRunCluster(t, 3)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
h := cmd.Handler.(*http.Handler).Handler
|
||||
|
||||
t.Run("PostSchema", func(t *testing.T) {
|
||||
w := httptest.NewRecorder()
|
||||
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/schema", strings.NewReader(`{"indexes":[{"name":"blah","options":{"keys":false,"trackExistence":true},"fields":[{"name":"f1","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false}}],"shardWidth":1048576}]}`)))
|
||||
if w.Code != gohttp.StatusNoContent {
|
||||
bod, err := ioutil.ReadAll(w.Result().Body)
|
||||
if err != nil {
|
||||
t.Errorf("reading body: %v", err)
|
||||
}
|
||||
t.Fatalf("unexpected code: %v, bod: %s", w.Code, bod)
|
||||
}
|
||||
for i := 0; i < len(cluster); i++ {
|
||||
cmd = cluster[i]
|
||||
idx, err := cmd.API.Index(context.Background(), "blah")
|
||||
if err != nil {
|
||||
t.Fatalf("getting index: %v", err)
|
||||
}
|
||||
if idx.Name() != "blah" {
|
||||
t.Fatalf("index did not get set, got %v", idx.Name())
|
||||
}
|
||||
|
||||
fld, err := cmd.API.Field(context.Background(), "blah", "f1")
|
||||
if err != nil {
|
||||
t.Fatalf("getting field: %v", err)
|
||||
}
|
||||
if fld.Name() != "f1" {
|
||||
t.Fatalf("unexpected field: %v", fld.Name())
|
||||
}
|
||||
}
|
||||
|
||||
h.ServeHTTP(w, test.MustNewHTTPRequest("DELETE", "/index/blah", nil))
|
||||
})
|
||||
}
|
||||
|
||||
func TestHandler_Endpoints(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
h := cmd.Handler.(*http.Handler).Handler
|
||||
holder := cmd.Server.Holder()
|
||||
hldr := test.Holder{Holder: holder}
|
||||
|
||||
// Ensure the handler returns "not found" for invalid paths.
|
||||
t.Run("Not Found", func(t *testing.T) {
|
||||
w := httptest.NewRecorder()
|
||||
h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/no_such_path", nil))
|
||||
|
|
@ -51,6 +92,48 @@ func TestHandler_Endpoints(t *testing.T) {
|
|||
}
|
||||
})
|
||||
|
||||
t.Run("SchemaEmpty", func(t *testing.T) {
|
||||
w := httptest.NewRecorder()
|
||||
h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/schema", nil))
|
||||
if w.Code != gohttp.StatusOK {
|
||||
t.Fatalf("unexpected status code: %d", w.Code)
|
||||
}
|
||||
body := w.Body.String()
|
||||
if body != "{\"indexes\":null}\n" {
|
||||
t.Fatalf("unexpected empty schema: '%v'", body)
|
||||
}
|
||||
|
||||
})
|
||||
|
||||
t.Run("PostSchema", func(t *testing.T) {
|
||||
w := httptest.NewRecorder()
|
||||
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/schema", strings.NewReader(`{"indexes":[{"name":"blah","options":{"keys":false,"trackExistence":true},"fields":[{"name":"f1","options":{"type":"set","cacheType":"ranked","cacheSize":50000,"keys":false}}],"shardWidth":1048576}]}`)))
|
||||
if w.Code != gohttp.StatusNoContent {
|
||||
bod, err := ioutil.ReadAll(w.Result().Body)
|
||||
if err != nil {
|
||||
t.Errorf("reading body: %v", err)
|
||||
}
|
||||
t.Fatalf("unexpected code: %v, bod: %s", w.Code, bod)
|
||||
}
|
||||
idx, err := cmd.API.Index(context.Background(), "blah")
|
||||
if err != nil {
|
||||
t.Fatalf("getting index: %v", err)
|
||||
}
|
||||
if idx.Name() != "blah" {
|
||||
t.Fatalf("index did not get set, got %v", idx.Name())
|
||||
}
|
||||
|
||||
fld, err := cmd.API.Field(context.Background(), "blah", "f1")
|
||||
if err != nil {
|
||||
t.Fatalf("getting field: %v", err)
|
||||
}
|
||||
if fld.Name() != "f1" {
|
||||
t.Fatalf("unexpected field: %v", fld.Name())
|
||||
}
|
||||
|
||||
h.ServeHTTP(w, test.MustNewHTTPRequest("DELETE", "/index/blah", nil))
|
||||
})
|
||||
|
||||
t.Run("Info", func(t *testing.T) {
|
||||
w := httptest.NewRecorder()
|
||||
h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/info", nil))
|
||||
|
|
@ -660,6 +743,7 @@ func TestHandler_Endpoints(t *testing.T) {
|
|||
}
|
||||
|
||||
clus := test.MustRunCluster(t, 1, []server.CommandOption{test.OptAllowedOrigins([]string{"http://test/"})})
|
||||
defer clus.Close()
|
||||
w = httptest.NewRecorder()
|
||||
h := clus[0].Handler.(*http.Handler).Handler
|
||||
h.ServeHTTP(w, req)
|
||||
|
|
|
|||
|
|
@ -321,6 +321,7 @@ func TestConfig_Parse_DataDir(t *testing.T) {
|
|||
func TestMain_RecalculateHashes(t *testing.T) {
|
||||
const clusterSize = 5
|
||||
cluster := test.MustRunCluster(t, clusterSize)
|
||||
defer cluster.Close()
|
||||
|
||||
// Create the schema.
|
||||
client0 := cluster[0].Client()
|
||||
|
|
|
|||
|
|
@ -211,7 +211,9 @@ func TestStatsCount_SetProfileAttrs(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestStatsCount_APICalls(t *testing.T) {
|
||||
cmd := test.MustRunCluster(t, 1)[0]
|
||||
cluster := test.MustRunCluster(t, 1)
|
||||
defer cluster.Close()
|
||||
cmd := cluster[0]
|
||||
h := cmd.Handler.(*http.Handler).Handler
|
||||
holder := cmd.Server.Holder()
|
||||
hldr := test.Holder{Holder: holder}
|
||||
|
|
|
|||
|
|
@ -28,6 +28,7 @@ import (
|
|||
func TestNewCluster(t *testing.T) {
|
||||
numNodes := 3
|
||||
cluster := test.MustRunCluster(t, numNodes)
|
||||
defer cluster.Close()
|
||||
|
||||
coordinator := getCoordinator(cluster[0])
|
||||
for i := 1; i < numNodes; i++ {
|
||||
|
|
|
|||
6
view.go
6
view.go
|
|
@ -88,12 +88,14 @@ func (v *view) open() error {
|
|||
|
||||
if err := func() error {
|
||||
// Ensure the view's path exists.
|
||||
v.logger.Debugf("ensure view path exists: %s", v.path)
|
||||
if err := os.MkdirAll(v.path, 0777); err != nil {
|
||||
return errors.Wrap(err, "creating view directory")
|
||||
} else if err := os.MkdirAll(filepath.Join(v.path, "fragments"), 0777); err != nil {
|
||||
return errors.Wrap(err, "creating fragments directory")
|
||||
}
|
||||
|
||||
v.logger.Debugf("open fragments for index/field/view: %s/%s/%s", v.index, v.field, v.name)
|
||||
if err := v.openFragments(); err != nil {
|
||||
return errors.Wrap(err, "opening fragments")
|
||||
}
|
||||
|
|
@ -104,6 +106,7 @@ func (v *view) open() error {
|
|||
return err
|
||||
}
|
||||
|
||||
v.logger.Debugf("successfully opened index/field/view: %s/%s/%s", v.index, v.field, v.name)
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -130,14 +133,17 @@ func (v *view) openFragments() error {
|
|||
// Parse filename into integer.
|
||||
shard, err := strconv.ParseUint(filepath.Base(fi.Name()), 10, 64)
|
||||
if err != nil {
|
||||
v.logger.Debugf("WARNING: couldn't use non-integer file as shard in index/field/view %s/%s/%s: %s", v.index, v.field, v.name, fi.Name())
|
||||
continue
|
||||
}
|
||||
|
||||
v.logger.Debugf("open index/field/view/fragment: %s/%s/%s/%d", v.index, v.field, v.name, shard)
|
||||
frag := v.newFragment(v.fragmentPath(shard), shard)
|
||||
if err := frag.Open(); err != nil {
|
||||
return fmt.Errorf("open fragment: shard=%d, err=%s", frag.shard, err)
|
||||
}
|
||||
frag.RowAttrStore = v.rowAttrStore
|
||||
v.logger.Debugf("add index/field/view/fragment to view.fragments: %s/%s/%s/%d", v.index, v.field, v.name, shard)
|
||||
v.fragments[frag.shard] = frag
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue