diff --git a/api.go b/api.go index 1405b219e..8b30053fb 100644 --- a/api.go +++ b/api.go @@ -655,7 +655,30 @@ func (api *API) Schema(ctx context.Context) []*IndexInfo { return api.holder.limitedSchema() } -func (api *API) ApplySchema(ctx context.Context, s *Schema) error { +// 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) } @@ -1277,6 +1300,7 @@ const ( //apiStatsWithTags // not implemented //apiVersion // not implemented apiViews + apiApplySchema ) var methodsCommon = map[apiMethod]struct{}{ @@ -1310,4 +1334,5 @@ var methodsNormal = map[apiMethod]struct{}{ apiRemoveNode: {}, apiShardNodes: {}, apiViews: {}, + apiApplySchema: {}, } diff --git a/apimethod_string.go b/apimethod_string.go index cafeab5af..309648a62 100644 --- a/apimethod_string.go +++ b/apimethod_string.go @@ -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) { diff --git a/client.go b/client.go index d537705c8..3b815a6bd 100644 --- a/client.go +++ b/client.go @@ -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 } diff --git a/ctl/export_test.go b/ctl/export_test.go index 6cb611460..e9f7db8a7 100644 --- a/ctl/export_test.go +++ b/ctl/export_test.go @@ -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) diff --git a/ctl/import_test.go b/ctl/import_test.go index 01abf372c..dd4d4a9e6 100644 --- a/ctl/import_test.go +++ b/ctl/import_test.go @@ -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(""))) diff --git a/executor_test.go b/executor_test.go index c49938637..65440c03a 100644 --- a/executor_test.go +++ b/executor_test.go @@ -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}) diff --git a/go.mod b/go.mod index 8f412ab0c..513390d90 100644 --- a/go.mod +++ b/go.mod @@ -26,7 +26,12 @@ require ( github.com/spf13/viper v1.3.1 github.com/uber/jaeger-client-go v2.15.0+incompatible github.com/uber/jaeger-lib v1.5.0 - golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4 + 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 ) diff --git a/go.sum b/go.sum index d33a2f661..9bb418b94 100644 --- a/go.sum +++ b/go.sum @@ -105,15 +105,33 @@ github.com/xordataexchange/crypt v0.0.3-0.20170626215501-b2862e3d0a77/go.mod h1: 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= +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= diff --git a/http/client.go b/http/client.go index 4f2f1c208..74706588d 100644 --- a/http/client.go +++ b/http/client.go @@ -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") diff --git a/http/client_test.go b/http/client_test.go index 0bda138ca..7fa398efa 100644 --- a/http/client_test.go +++ b/http/client_test.go @@ -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} diff --git a/http/handler.go b/http/handler.go index c6335e3dd..b52ee52a3 100644 --- a/http/handler.go +++ b/http/handler.go @@ -187,7 +187,7 @@ func (h *Handler) populateValidators() { h.validators["GetInfo"] = queryValidationSpecRequired() h.validators["RecalculateCaches"] = queryValidationSpecRequired() h.validators["GetSchema"] = queryValidationSpecRequired() - h.validators["PostSchema"] = queryValidationSpecRequired() + h.validators["PostSchema"] = queryValidationSpecRequired().Optional("remote") h.validators["GetStatus"] = queryValidationSpecRequired() h.validators["GetVersion"] = queryValidationSpecRequired() h.validators["PostClusterMessage"] = queryValidationSpecRequired() @@ -420,13 +420,20 @@ func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) { } 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); err != nil { + 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 } diff --git a/http/translator_test.go b/http/translator_test.go index bec2089a5..87bb4366d 100644 --- a/http/translator_test.go +++ b/http/translator_test.go @@ -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) diff --git a/server/cluster_test.go b/server/cluster_test.go index 04ff8ef68..b6f6cf6bf 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -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] diff --git a/server/config_internal_test.go b/server/config_internal_test.go index 2b29ea51e..ca7548895 100644 --- a/server/config_internal_test.go +++ b/server/config_internal_test.go @@ -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) + } + }) } } diff --git a/server/handler_test.go b/server/handler_test.go index 1f56aa6f3..fc832269a 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -36,8 +36,49 @@ import ( "github.com/pilosa/pilosa/test" ) +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} @@ -68,7 +109,11 @@ func TestHandler_Endpoints(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 { - t.Fatalf("unexpected code: %v", w.Code) + 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 { @@ -698,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) diff --git a/server/server_test.go b/server/server_test.go index e911b25e6..11a026dea 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -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() diff --git a/stats/stats_test.go b/stats/stats_test.go index 3da83ce70..f09cce79e 100644 --- a/stats/stats_test.go +++ b/stats/stats_test.go @@ -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} diff --git a/test/pilosa_test.go b/test/pilosa_test.go index 1f00ff68b..1d0feb237 100644 --- a/test/pilosa_test.go +++ b/test/pilosa_test.go @@ -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++ {