Merge branch 'develop' into unexport-api-translatestore

This commit is contained in:
Cody Soyland 2018-06-27 13:36:19 -05:00
commit c042c77499
15 changed files with 323 additions and 213 deletions

6
Gopkg.lock generated
View file

@ -197,12 +197,6 @@
revision = "645ef00459ed84a119197bfb8d8205042c6df63d"
version = "v0.8.0"
[[projects]]
name = "github.com/rakyll/statik"
packages = ["fs"]
revision = "fd36b3595eb2ec8da4b8153b107f7ea08504899d"
version = "v0.1.1"
[[projects]]
name = "github.com/satori/go.uuid"
packages = ["."]

View file

@ -27,7 +27,7 @@ import (
type MemberSet interface {
// Open starts any network activity implemented by the MemberSet
// Node is the local node, used for membership broadcasts.
Open(n *Node) error
Open() error
}
// StaticMemberSet represents a basic MemberSet for testing.
@ -43,7 +43,7 @@ func NewStaticMemberSet(nodes []*Node) *StaticMemberSet {
}
// Open implements the MemberSet interface to start network activity, but for a static MemberSet it does nothing.
func (s *StaticMemberSet) Open(n *Node) error {
func (s *StaticMemberSet) Open() error {
return nil
}

View file

@ -56,6 +56,7 @@ func testMessageMarshal(t *testing.T, m proto.Message) {
// Ensure that BroadcastReceiver can register a BroadcastHandler.
func TestBroadcast_BroadcastReceiver(t *testing.T) {
t.Skip("broadcast receiver")
path, err := ioutil.TempDir("", "pilosa-")
if err != nil {
panic(err)
@ -67,24 +68,24 @@ func TestBroadcast_BroadcastReceiver(t *testing.T) {
if err != nil {
t.Fatalf("setting up server: %v", err)
}
s := com.Server
// s := com.Server
sbr := NewSimpleBroadcastReceiver()
sbh := NewSimpleBroadcastHandler()
// sbr := NewSimpleBroadcastReceiver()
// sbh := NewSimpleBroadcastHandler()
s.BroadcastReceiver = sbr
s.BroadcastReceiver.Start(sbh)
// s.BroadcastReceiver = sbr
// s.BroadcastReceiver.Start(sbh)
msg := &internal.DeleteIndexMessage{
Index: "i",
}
// msg := &internal.DeleteIndexMessage{
// Index: "i",
// }
s.BroadcastReceiver.(*SimpleBroadcastReceiver).Receive(msg)
// s.BroadcastReceiver.(*SimpleBroadcastReceiver).Receive(msg)
// Make sure the message received is what was sentd
if !reflect.DeepEqual(sbh.receivedMessage, msg) {
t.Fatalf("unexpected message: %s", sbh.receivedMessage)
}
// // Make sure the message received is what was sentd
// if !reflect.DeepEqual(sbh.receivedMessage, msg) {
// t.Fatalf("unexpected message: %s", sbh.receivedMessage)
// }
}
type SimpleBroadcastReceiver struct {

View file

@ -861,7 +861,7 @@ func (h *jmphasher) Hash(key uint64, n int) int {
return int(b)
}
func (c *Cluster) open() error {
func (c *Cluster) setup() error {
// Cluster always comes up in state STARTING until cluster membership is determined.
c.state = ClusterStateStarting
@ -876,7 +876,7 @@ func (c *Cluster) open() error {
if c.isCoordinator() {
err := c.considerTopology()
if err != nil {
return fmt.Errorf("considerTopology: %v", err)
return errors.Wrap(err, "considerTopology")
}
}
@ -885,15 +885,21 @@ func (c *Cluster) open() error {
if err != nil {
return errors.Wrap(err, "adding local node")
}
return nil
}
// Start the EventReceiver.
if err := c.EventReceiver.Start(c); err != nil {
return fmt.Errorf("starting EventReceiver: %v", err)
func (c *Cluster) open() error {
err := c.setup()
if err != nil {
return errors.Wrap(err, "setting up cluster")
}
return c.waitForStarted()
}
func (c *Cluster) waitForStarted() error {
// Open MemberSet communication.
if err := c.MemberSet.Open(c.Node); err != nil {
return fmt.Errorf("opening MemberSet: %v", err)
if err := c.MemberSet.Open(); err != nil {
return errors.Wrap(err, "opening MemberSet")
}
// If not coordinator then wait for ClusterStatus from coordinator.

View file

@ -25,6 +25,7 @@ import (
"github.com/davecgh/go-spew/spew"
"github.com/pilosa/pilosa/internal"
"github.com/pkg/errors"
)
// Ensure that fragCombos creates the correct fragment mapping.
@ -591,10 +592,10 @@ func TestCluster_ResizeStates(t *testing.T) {
tc.WriteTopology(node.Path, top)
// Open TestCluster.
expected := "considerTopology: coordinator node0 is not in topology: [some-other-host]"
expected := "coordinator node0 is not in topology: [some-other-host]"
err := tc.Open()
if err == nil || err.Error() != expected {
t.Errorf("did not receive expected error: %s", expected)
if err == nil || errors.Cause(err).Error() != expected {
t.Errorf("did not receive expected error, got: %s", errors.Cause(err).Error())
}
// Close TestCluster.

View file

@ -46,16 +46,16 @@ func main() {
panic(err)
}
// We need to refer to indexes and frames before we can use them in a query.
// We need to refer to indexes and fields before we can use them in a query.
repository, _ := schema.Index("repository")
stargazer, _ := repository.Frame("stargazer")
language, _ := repository.Frame("language")
stargazer, _ := repository.Field("stargazer")
language, _ := repository.Field("language")
var response *pilosa.QueryResponse
// Which repositories did user 14 star:
response, _ = client.Query(stargazer.Bitmap(14))
fmt.Println("User 14 starred: ", response.Result().Bitmap().Bits)
response, _ = client.Query(stargazer.Row(14))
fmt.Println("User 14 starred: ", response.Result().Row().Columns)
// What are the top 5 languages in the sample data?
response, err = client.Query(language.TopN(5))
@ -68,26 +68,26 @@ func main() {
// Which repositories were starred by both user 14 and 19:
response, _ = client.Query(
repository.Intersect(
stargazer.Bitmap(14),
stargazer.Bitmap(19)))
fmt.Println("Both user 14 and 19 starred:", response.Result().Bitmap().Bits)
stargazer.Row(14),
stargazer.Row(19)))
fmt.Println("Both user 14 and 19 starred:", response.Result().Row().Columns)
// Which repositories were starred by user 14 or 19:
response, _ = client.Query(
repository.Union(
stargazer.Bitmap(14),
stargazer.Bitmap(19)))
fmt.Println("User 14 or 19 starred:", response.Result().Bitmap().Bits)
stargazer.Row(14),
stargazer.Row(19)))
fmt.Println("User 14 or 19 starred:", response.Result().Row().Columns)
// Which repositories were starred by user 14 or 19 and were written in language 1:
response, _ = client.Query(
repository.Intersect(
repository.Union(
stargazer.Bitmap(14),
stargazer.Bitmap(19),
stargazer.Row(14),
stargazer.Row(19),
),
language.Bitmap(1)))
fmt.Println("User 14 or 19 starred, written in language 1:", response.Result().Bitmap().Bits)
language.Row(1)))
fmt.Println("User 14 or 19 starred, written in language 1:", response.Result().Row().Columns)
// Set user 99999 as a stargazer for repository 77777?
client.Query(stargazer.SetBit(99999, 77777))
@ -112,6 +112,7 @@ We are going to use the index you have created in the [Getting Started](../getti
Error handling has been omitted in the example below for brevity.
```python
from __future__ import print_function
from pilosa import Index, Client, PilosaError, TimeQuantum
# We will just use the default client which assumes the server is at http://localhost:10101
@ -122,8 +123,8 @@ client = Client()
# and the stargazer data should be imported.
# See the Getting Started repository: https://github.com/pilosa/getting-started/
# Let's create Index and Frame objects, which will contain the settings
# for the corresponding indexes and frames.
# Let's create Index and Field objects, which will contain the settings
# for the corresponding indexes and fields.
try:
schema = client.schema()
except PilosaError as e:
@ -132,13 +133,13 @@ except PilosaError as e:
# We will just terminate the program in this case.
raise SystemExit(e)
# We need to refer to indexes and frames before we can use them in a query.
# We need to refer to indexes and fields before we can use them in a query.
repository = schema.index("repository")
stargazer = repository.frame("stargazer")
language = repository.frame("language")
stargazer = repository.field("stargazer")
language = repository.field("language")
# Which repositories did user 8 star:
repository_ids = client.query(stargazer.bitmap(14)).result.bitmap.bits
repository_ids = client.query(stargazer.row(14)).result.row.columns
print("User 8 starred: ", repository_ids)
# What are the top 5 languages in the sample data:
@ -147,29 +148,29 @@ print("Top 5 languages: ", [item.id for item in top_languages])
# Which repositories were starred by both user 14 and 19:
query = repository.intersect(
stargazer.bitmap(14),
stargazer.bitmap(19)
stargazer.row(14),
stargazer.row(19)
)
mutually_starred = client.query(query).result.bitmap.bits
mutually_starred = client.query(query).result.row.columns
print("Both user 14 and 19 starred:", mutually_starred)
# Which repositories were starred by user 14 or 19:
query = repository.union(
stargazer.bitmap(14),
stargazer.bitmap(19)
stargazer.row(14),
stargazer.row(19)
)
either_starred = client.query(query).result.bitmap.bits
either_starred = client.query(query).result.row.columns
print("User 14 or 19 starred:", either_starred)
# Which repositories were starred by user 14 or 19 and were written in language 1:
query = repository.intersect(
repository.union(
stargazer.bitmap(14),
stargazer.bitmap(19)
stargazer.row(14),
stargazer.row(19)
),
language.bitmap(1)
language.row(1)
)
mutually_starred = client.query(query).result.bitmap.bits
mutually_starred = client.query(query).result.row.columns
print("User 14 or 19 starred, written in language 1:", mutually_starred)
# Set user 99999 as a stargazer for repository 77777
@ -218,10 +219,10 @@ public class StarTrace {
throw new RuntimeException(ex);
}
// We need to refer to indexes and frames before we can use them in a query.
// We need to refer to indexes and fields before we can use them in a query.
Index repository = schema.index("repository");
Frame stargazer = repository.frame("stargazer");
Frame language = repository.frame("language");
Field stargazer = repository.field("stargazer");
Field language = repository.field("language");
QueryResponse response;
QueryResult result;
@ -229,8 +230,8 @@ public class StarTrace {
List<Long> repositoryIDs;
// Which repositories did user 14 star:
response = client.query(stargazer.bitmap(14));
repositoryIDs = response.getResult().getBitmap().getBits();
response = client.query(stargazer.row(14));
repositoryIDs = response.getResult().getRow().getColumns();
System.out.println("User 14 starred: " + repositoryIDs);
// What are the top 5 languages in the sample data:
@ -245,32 +246,32 @@ public class StarTrace {
// Which repositories were starred by both user 14 and 19:
query = repository.intersect(
stargazer.bitmap(14),
stargazer.bitmap(19)
stargazer.row(14),
stargazer.row(19)
);
response = client.query(query);
repositoryIDs = response.getResult().getBitmap().getBits();
repositoryIDs = response.getResult().getRow().getColumns();
System.out.println("Both user 14 and 19 starred: " + repositoryIDs);
// Which repositories were starred by user 14 or 19:
query = repository.union(
stargazer.bitmap(14),
stargazer.bitmap(19)
stargazer.row(14),
stargazer.row(19)
);
response = client.query(query);
repositoryIDs = response.getResult().getBitmap().getBits();
repositoryIDs = response.getResult().getRow().getColumns();
System.out.println("User 14 or 19 starred: " + repositoryIDs);
// Which repositories were starred by user 14 or 19 and were written in language 1:
query = repository.intersect(
repository.union(
stargazer.bitmap(14),
stargazer.bitmap(19)
stargazer.row(14),
stargazer.row(19)
),
language.bitmap(1)
language.row(1)
);
response = client.query(query);
repositoryIDs = response.getResult().getBitmap().getBits();
repositoryIDs = response.getResult().getRow().getColumns();
System.out.println("User 14 or 19 starred, written in language 1: " + repositoryIDs);
// Set user 99999 as a stargazer for repository 77777:

View file

@ -34,14 +34,15 @@ Let's make sure Pilosa is running:
curl localhost:10101/status
```
``` response
{"state":"NORMAL","nodes":[{"id":"18eb5546-5a1a-4ba4-9c52-b53fbe22317e","uri":{"scheme":"http","host":"localhost","port":10101}}]}
{"state":"NORMAL","nodes":[{"id":"91715a50-7d50-4c54-9a03-873801da1cd1","uri":{"scheme":"http","host":"localhost","port
":10101},"isCoordinator":true}],"localID":"91715a50-7d50-4c54-9a03-873801da1cd1"}
```
### Sample Project
In order to better understand Pilosa's capabilities, we will create a sample project called "Star Trace" containing information about 1,000 popular Github repositories which have "go" in their name. The Star Trace index will include data points such as programming language, tags, and stargazers—people who have starred a project.
Although Pilosa doesn't keep the data in a tabular format, we still use the terms "columns" and "rows" when describing the data model. We put the primary objects in columns, and the properties of those objects in rows. For example, the Star Trace project will contain an index called "repository" which contains columns representing Github repositories, and rows representing properties like programming languages and tags. We can better organize the rows by grouping them into sets called Frames. So the "repository" index might have a "languages" frame as well as a "tags" frame. You can learn more about indexes and frames in the [Data Model](../data-model/) section of the documentation.
Although Pilosa doesn't keep the data in a tabular format, we still use the terms "columns" and "rows" when describing the data model. We put the primary objects in columns, and the properties of those objects in rows. For example, the Star Trace project will contain an index called "repository" which contains columns representing Github repositories, and rows representing properties like programming languages and tags. We can better organize the rows by grouping them into sets called Fields. So the "repository" index might have a "languages" field as well as a "tags" field. You can learn more about indexes and fields in the [Data Model](../data-model/) section of the documentation.
#### Create the Schema
@ -55,7 +56,7 @@ curl localhost:10101/schema
{"indexes":null}
```
Before we can import data or run queries, we need to create our indexes and the frames within them. Let's create the repository index first:
Before we can import data or run queries, we need to create our indexes and the fields within them. Let's create the repository index first:
``` request
curl localhost:10101/index/repository -X POST
```
@ -63,27 +64,29 @@ curl localhost:10101/index/repository -X POST
{}
```
Let's create the `stargazer` frame which has user IDs of stargazers as its rows:
Let's create the `stargazer` field which has user IDs of stargazers as its rows:
``` request
curl localhost:10101/index/repository/frame/stargazer \
curl localhost:10101/index/repository/field/stargazer \
-X POST \
-d '{"options": {"timeQuantum": "YMD"}}'
-d '{"options": {"type": "time", "timeQuantum": "YMD"}}'
```
``` response
{}
```
Since our data contains time stamps for the time users starred repos, we set the *time quantum* for the `stargazer` frame in the options as well. 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 for 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` frame, which will contain IDs for programming languages:
Next up is the `language` field, which will contain IDs for programming languages:
``` request
curl localhost:10101/index/repository/frame/language \
curl localhost:10101/index/repository/field/language \
-X POST
```
``` response
{}
```
The `language` is a `set` field, but since the default field type is `set`, we didn't specify it in field options.
#### Import Data From CSV Files
Download the `stargazer.csv` and `language.csv` files here:
@ -116,14 +119,14 @@ Which repositories did user 14 star:
``` request
curl localhost:10101/index/repository/query \
-X POST \
-d 'Bitmap(frame="stargazer", row=14)'
-d 'Bitmap(field="stargazer", row=14)'
```
``` response
{
"results":[
{
"attrs":{},
"bits":[1,2,3,362,368,391,396,409,416,430,436,450,454,460,461,464,466,469,470,483,484,486,490,491,503,504,514]
"columns":[1,2,3,362,368,391,396,409,416,430,436,450,454,460,461,464,466,469,470,483,484,486,490,491,503,504,514]
}
]
}
@ -133,7 +136,7 @@ What are the top 5 languages in the sample data:
``` request
curl localhost:10101/index/repository/query \
-X POST \
-d 'TopN(frame="language", n=5)'
-d 'TopN(field="language", n=5)'
```
``` response
{
@ -154,8 +157,8 @@ Which repositories were starred by user 14 and 19:
curl localhost:10101/index/repository/query \
-X POST \
-d 'Intersect(
Bitmap(frame="stargazer", row=14),
Bitmap(frame="stargazer", row=19)
Bitmap(field="stargazer", row=14),
Bitmap(field="stargazer", row=19)
)'
```
``` response
@ -163,7 +166,7 @@ curl localhost:10101/index/repository/query \
"results":[
{
"attrs":{},
"bits":[2,3,362,396,416,461,464,466,470,486]
"columns":[2,3,362,396,416,461,464,466,470,486]
}
]
}
@ -174,8 +177,8 @@ Which repositories were starred by user 14 or 19:
curl localhost:10101/index/repository/query \
-X POST \
-d 'Union(
Bitmap(frame="stargazer", row=14),
Bitmap(frame="stargazer", row=19)
Bitmap(field="stargazer", row=14),
Bitmap(field="stargazer", row=19)
)'
```
``` response
@ -183,7 +186,7 @@ curl localhost:10101/index/repository/query \
"results":[
{
"attrs":{},
"bits":[1,2,3,361,362,368,376,377,378,382,386,388,391,396,398,400,409,411,412,416,426,428,430,435,436,450,452,453,454,456,460,461,464,465,466,469,470,483,484,486,487,489,490,491,500,503,504,505,512,514]
"columns":[1,2,3,361,362,368,376,377,378,382,386,388,391,396,398,400,409,411,412,416,426,428,430,435,436,450,452,453,454,456,460,461,464,465,466,469,470,483,484,486,487,489,490,491,500,503,504,505,512,514]
}
]
}
@ -194,9 +197,9 @@ Which repositories were starred by user 14 and 19 and also were written in langu
curl localhost:10101/index/repository/query \
-X POST \
-d 'Intersect(
Bitmap(frame="stargazer", row=14),
Bitmap(frame="stargazer", row=19),
Bitmap(frame="language", row=1)
Bitmap(field="stargazer", row=14),
Bitmap(field="stargazer", row=19),
Bitmap(field="language", row=1)
)'
```
``` response
@ -204,7 +207,7 @@ curl localhost:10101/index/repository/query \
"results":[
{
"attrs":{},
"bits":[2,362,416,461]
"columns":[2,362,416,461]
}
]
}
@ -214,7 +217,7 @@ Set user 99999 as a stargazer for repository 77777:
``` request
curl localhost:10101/index/repository/query \
-X POST \
-d 'SetBit(frame="stargazer", col=77777, row=99999)'
-d 'SetBit(field="stargazer", col=77777, row=99999)'
```
``` response
{"results":[true]}

View file

@ -1606,6 +1606,9 @@ func (e *Executor) translateCall(index string, idx *Index, c *pql.Call) error {
}
// Translate column key.
if idx.Keys() {
if c.Args[colKey] != nil && !isString(c.Args[colKey]) {
return errors.New("column value must be a string when index 'keys' option enabled")
}
if value := callArgString(c, colKey); value != "" {
ids, err := e.TranslateStore.TranslateColumnsToUint64(index, []string{value})
if err != nil {
@ -1613,6 +1616,10 @@ func (e *Executor) translateCall(index string, idx *Index, c *pql.Call) error {
}
c.Args[colKey] = ids[0]
}
} else {
if isString(c.Args[colKey]) {
return errors.New("string 'col' value not allowed unless index 'keys' option enabled")
}
}
// Translate row key, if field is specified & key exists.
@ -1622,6 +1629,9 @@ func (e *Executor) translateCall(index string, idx *Index, c *pql.Call) error {
return ErrFieldNotFound
}
if field.Keys() {
if c.Args[rowKey] != nil && !isString(c.Args[rowKey]) {
return errors.New("row value must be a string when field 'keys' option enabled")
}
if value := callArgString(c, rowKey); value != "" {
ids, err := e.TranslateStore.TranslateRowsToUint64(index, fieldName, []string{value})
if err != nil {
@ -1629,6 +1639,10 @@ func (e *Executor) translateCall(index string, idx *Index, c *pql.Call) error {
}
c.Args[rowKey] = ids[0]
}
} else {
if isString(c.Args[rowKey]) {
return errors.New("string 'row' value not allowed unless field 'keys' option enabled")
}
}
}
@ -1801,3 +1815,8 @@ func callArgString(call *pql.Call, key string) string {
s, _ := value.(string)
return s
}
func isString(v interface{}) bool {
_, ok := v.(string)
return ok
}

View file

@ -27,6 +27,7 @@ import (
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/pql"
"github.com/pilosa/pilosa/test"
"github.com/pkg/errors"
)
// Ensure a bitmap query can be executed.
@ -264,36 +265,107 @@ func TestExecutor_Execute_Count(t *testing.T) {
}
// Ensure a set query can be executed.
func TestExecutor_Execute_Set(t *testing.T) {
hldr := test.MustOpenHolder()
defer hldr.Close()
func TestExecutor_Execute_SetBit(t *testing.T) {
t.Run("ID", func(t *testing.T) {
cmd := test.MustRunMainWithCluster(t, 1)[0]
holder := cmd.Server.Holder()
hldr := test.Holder{Holder: holder}
hldr.SetBit("i", "f", 1, 0)
// set a bit so the view gets created.
hldr.SetBit("i", "f", 1, 0)
t.Run("OK", func(t *testing.T) {
hldr.ClearBit("i", "f", 11, 1)
if n := hldr.Row("i", "f", 11).Count(); n != 0 {
t.Fatalf("unexpected bitmap count: %d", n)
}
e := test.NewExecutor(hldr.Holder, pilosa.NewTestCluster(1))
if n := hldr.Row("i", "f", 11).Count(); n != 0 {
t.Fatalf("unexpected bitmap count: %d", n)
}
if res, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(1, f=11)`}); err != nil {
t.Fatal(err)
} else {
if !res.Results[0].(bool) {
t.Fatalf("expected column changed")
}
}
if res, err := e.Execute(context.Background(), "i", test.MustParse(`Set(1, f=11)`), nil, nil); err != nil {
t.Fatal(err)
} else {
if !res[0].(bool) {
t.Fatalf("expected column changed")
}
}
if n := hldr.Row("i", "f", 11).Count(); n != 1 {
t.Fatalf("unexpected bitmap count: %d", n)
}
if res, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(1, f=11)`}); err != nil {
t.Fatal(err)
} else {
if res.Results[0].(bool) {
t.Fatalf("expected column unchanged")
}
}
})
if n := hldr.Row("i", "f", 11).Count(); n != 1 {
t.Fatalf("unexpected bitmap count: %d", n)
}
if res, err := e.Execute(context.Background(), "i", test.MustParse(`Set(1, f=11)`), nil, nil); err != nil {
t.Fatal(err)
} else {
if res[0].(bool) {
t.Fatalf("expected column unchanged")
}
}
t.Run("ErrInvalidColValueType", func(t *testing.T) {
if _, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set("foo", f=1)`}); err == nil || errors.Cause(err).Error() != `string 'col' value not allowed unless index 'keys' option enabled` {
t.Fatalf("The error is: '%v'", err)
}
})
t.Run("ErrInvalidRowValueType", func(t *testing.T) {
if _, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(2, f="bar")`}); err == nil || errors.Cause(err).Error() != `string 'row' value not allowed unless field 'keys' option enabled` {
t.Fatal(err)
}
})
})
t.Run("Keys", func(t *testing.T) {
cmd := test.MustRunMainWithCluster(t, 1)[0]
holder := cmd.Server.Holder()
hldr := test.Holder{Holder: holder}
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{Keys: true})
t.Run("OK", func(t *testing.T) {
hldr.SetBit("i", "f", 1, 0)
if n := hldr.Row("i", "f", 11).Count(); n != 0 {
t.Fatalf("unexpected bitmap count: %d", n)
}
if res, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set("foo", f=11)`}); err != nil {
t.Fatal(err)
} else {
if !res.Results[0].(bool) {
t.Fatalf("expected column changed")
}
}
if n := hldr.Row("i", "f", 11).Count(); n != 1 {
t.Fatalf("unexpected bitmap count: %d", n)
}
if res, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set("foo", f=11)`}); err != nil {
t.Fatal(err)
} else {
if res.Results[0].(bool) {
t.Fatalf("expected column unchanged")
}
}
})
t.Run("ErrInvalidColValueType", func(t *testing.T) {
if err := index.DeleteField("f"); err != nil {
t.Fatal(err)
}
if _, err := index.CreateField("f", pilosa.FieldOptions{}); err != nil {
t.Fatal(err)
}
if _, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(2, f=1)`}); err == nil || errors.Cause(err).Error() != `column value must be a string when index 'keys' option enabled` {
t.Fatal(err)
}
})
t.Run("ErrInvalidRowValueType", func(t *testing.T) {
index := hldr.MustCreateIndexIfNotExists("inokey", pilosa.IndexOptions{})
if _, err := index.CreateField("f", pilosa.FieldOptions{Keys: true}); err != nil {
t.Fatal(err)
}
if _, err := cmd.API.Query(context.Background(), &pilosa.QueryRequest{Index: "inokey", Query: `Set(2, f=1)`}); err == nil || errors.Cause(err).Error() != `row value must be a string when field 'keys' option enabled` {
t.Fatal(err)
}
})
})
}
// Ensure old PQL syntax doesn't break anything too badly.

View file

@ -33,31 +33,25 @@ import (
)
// Ensure GossipMemberSet implements interfaces.
var _ pilosa.BroadcastReceiver = &GossipMemberSet{}
var _ memberlist.Delegate = &GossipMemberSet{}
// GossipMemberSet represents a gossip implementation of MemberSet using memberlist.
type GossipMemberSet struct {
mu sync.RWMutex
node *pilosa.Node
memberlist *memberlist.Memberlist
handler pilosa.BroadcastHandler
broadcasts *memberlist.TransmitLimitedQueue
statusHandler pilosa.StatusHandler
config *gossipConfig
pserver *pilosa.Server
config *gossipConfig
Logger pilosa.Logger
logger *log.Logger
transport *Transport
}
// Start implements the BroadcastReceiver interface and sets the BroadcastHandler.
func (g *GossipMemberSet) Start(h pilosa.BroadcastHandler) error {
g.handler = h
return nil
gossipEventReceiver *GossipEventReceiver
}
// GetBindAddr returns the gossip bind address based on config and auto bind port.
@ -68,14 +62,15 @@ func (g *GossipMemberSet) GetBindAddr() string {
}
// Open implements the MemberSet interface to start network activity.
func (g *GossipMemberSet) Open(n *pilosa.Node) error {
func (g *GossipMemberSet) Open() error {
err := g.gossipEventReceiver.Start(g.pserver)
if err != nil {
return errors.Wrap(err, "starting event delegate")
}
if g.handler == nil {
return fmt.Errorf("must call Start(pilosa.BroadcastHandler) before calling Open()")
}
g.node = n
err := error(nil)
g.mu.Lock()
g.memberlist, err = memberlist.Create(g.config.memberlistConfig)
g.mu.Unlock()
@ -166,7 +161,7 @@ func WithLogger(logger *log.Logger) GossipMemberSetOption {
}
// NewGossipMemberSet returns a new instance of GossipMemberSet based on options.
func NewGossipMemberSet(name string, host string, cfg Config, ger *GossipEventReceiver, sh pilosa.StatusHandler, options ...GossipMemberSetOption) (*GossipMemberSet, error) {
func NewGossipMemberSet(name string, host string, cfg Config, s *pilosa.Server, options ...GossipMemberSetOption) (*GossipMemberSet, error) {
g := &GossipMemberSet{
Logger: pilosa.NopLogger,
}
@ -177,6 +172,10 @@ func NewGossipMemberSet(name string, host string, cfg Config, ger *GossipEventRe
return nil, errors.Wrap(err, "executing option")
}
}
ger := NewGossipEventReceiver(g.logger)
g.gossipEventReceiver = ger
g.handler = s
if g.transport == nil {
port, err := strconv.Atoi(cfg.Port)
@ -232,14 +231,14 @@ func NewGossipMemberSet(name string, host string, cfg Config, ger *GossipEventRe
gossipSeeds: cfg.Seeds,
}
g.statusHandler = sh
g.pserver = s
return g, nil
}
// NodeMeta implementation of the memberlist.Delegate interface.
func (g *GossipMemberSet) NodeMeta(limit int) []byte {
buf, err := proto.Marshal(pilosa.EncodeNode(g.node))
buf, err := proto.Marshal(pilosa.EncodeNode(g.pserver.Node()))
if err != nil {
g.Logger.Printf("marshal message error: %s", err)
return []byte{}
@ -270,7 +269,7 @@ func (g *GossipMemberSet) GetBroadcasts(overhead, limit int) [][]byte {
// LocalState implementation of the memberlist.Delegate interface
// sends this Node's state data.
func (g *GossipMemberSet) LocalState(join bool) []byte {
pb, err := g.statusHandler.LocalStatus()
pb, err := g.pserver.LocalStatus()
if err != nil {
g.Logger.Printf("error getting local state, err=%s", err)
return []byte{}
@ -294,7 +293,7 @@ func (g *GossipMemberSet) MergeRemoteState(buf []byte, join bool) {
g.Logger.Printf("error unmarshalling nodestate data, err=%s", err)
return
}
err := g.statusHandler.HandleRemoteStatus(&pb)
err := g.pserver.HandleRemoteStatus(&pb)
if err != nil {
g.Logger.Printf("merge state error: %s", err)
}
@ -309,11 +308,11 @@ type GossipEventReceiver struct {
ch chan memberlist.NodeEvent
eventHandler pilosa.EventHandler
Logger pilosa.Logger
Logger *log.Logger
}
// NewGossipEventReceiver returns a new instance of GossipEventReceiver.
func NewGossipEventReceiver(logger pilosa.Logger) *GossipEventReceiver {
func NewGossipEventReceiver(logger *log.Logger) *GossipEventReceiver {
return &GossipEventReceiver{
ch: make(chan memberlist.NodeEvent, 1),
Logger: logger,

View file

@ -4,6 +4,7 @@ import (
"context"
"io"
"io/ioutil"
gohttp "net/http"
"testing"
"time"
@ -15,8 +16,6 @@ import (
)
func TestTranslateStore_Reader(t *testing.T) {
t.Skip() // Until test.NewServer() works
// Ensure client can connect and stream the translate store data.
t.Run("OK", func(t *testing.T) {
t.Run("ServerDisconnect", func(t *testing.T) {
@ -46,15 +45,30 @@ func TestTranslateStore_Reader(t *testing.T) {
// Setup handler on test server.
var translateStore mock.TranslateStore
translateStore.ReaderFunc = func(ctx context.Context, off int64) (io.ReadCloser, error) {
if off != 100 {
t.Fatalf("unexpected off: %d", off)
// Check context to make sure this is the call we are looking for.
// (Something else calls ReaderFunc on server startup)
if ctx.Value(gohttp.ServerContextKey) != nil {
if off != 100 {
t.Fatalf("unexpected off: %d", off)
}
return &mrc, nil
}
return &mrc, nil
mrc2 := mock.ReadCloser{
ReadFunc: func(p []byte) (int, error) {
return 0, io.EOF
},
CloseFunc: func() error {
return nil
},
}
return &mrc2, nil
}
opts := server.OptCommandServerOptions(pilosa.OptServerPrimaryTranslateStore(translateStore))
main := test.MustRunMainWithCluster(t, 1, []server.CommandOption{opts})[0]
defer main.Close()
// Connect to server and stream all available data.
@ -128,6 +142,7 @@ func TestTranslateStore_Reader(t *testing.T) {
opts := server.OptCommandServerOptions(pilosa.OptServerPrimaryTranslateStore(translateStore))
main := test.MustRunMainWithCluster(t, 1, []server.CommandOption{opts})[0]
defer main.Close()
_, err := http.NewTranslateStore(main.Server.URI.String()).Reader(context.Background(), 0)
if err != pilosa.ErrNotImplemented {

View file

@ -61,10 +61,9 @@ type Server struct {
clusterDisabled bool
// External
BroadcastReceiver BroadcastReceiver
systemInfo SystemInfo
gcNotifier GCNotifier
logger Logger
systemInfo SystemInfo
gcNotifier GCNotifier
logger Logger
NodeID string
URI URI
@ -72,6 +71,7 @@ type Server struct {
metricInterval time.Duration
diagnosticInterval time.Duration
maxWritesPerRequest int
isCoordinator bool
primaryTranslateStore TranslateStore
@ -204,15 +204,21 @@ func OptServerClusterDisabled(disabled bool, hosts []string) ServerOption {
}
}
func OptServerIsCoordinator(is bool) ServerOption {
return func(s *Server) error {
s.isCoordinator = is
return nil
}
}
// NewServer returns a new instance of Server.
func NewServer(opts ...ServerOption) (*Server, error) {
s := &Server{
closing: make(chan struct{}),
Cluster: NewCluster(),
holder: NewHolder(),
BroadcastReceiver: NopBroadcastReceiver,
diagnostics: NewDiagnosticsCollector(DefaultDiagnosticServer),
systemInfo: NewNopSystemInfo(),
closing: make(chan struct{}),
Cluster: NewCluster(),
holder: NewHolder(),
diagnostics: NewDiagnosticsCollector(DefaultDiagnosticServer),
systemInfo: NewNopSystemInfo(),
gcNotifier: NopGCNotifier,
@ -246,14 +252,15 @@ func NewServer(opts ...ServerOption) (*Server, error) {
// Initialize translation database.
s.translateFile = NewTranslateFile()
s.translateFile.Path = filepath.Join(path, "keys")
s.translateFile.Path = filepath.Join(path, ".keys")
s.translateFile.PrimaryTranslateStore = s.primaryTranslateStore
if err := s.translateFile.Open(); err != nil {
return nil, err
}
// Get or create NodeID.
s.NodeID = s.LoadNodeID()
if s.isCoordinator {
s.Cluster.Coordinator = s.NodeID
}
// Set Cluster Node.
node := &Node{
ID: s.NodeID,
@ -276,6 +283,14 @@ func NewServer(opts ...ServerOption) (*Server, error) {
s.executor.Cluster = s.Cluster
s.executor.TranslateStore = s.translateFile
s.executor.MaxWritesPerRequest = s.maxWritesPerRequest
s.Cluster.Broadcaster = s
s.Cluster.MaxWritesPerRequest = s.maxWritesPerRequest
s.holder.Broadcaster = s
err = s.Cluster.setup()
if err != nil {
return nil, errors.Wrap(err, "setting up cluster")
}
return s, nil
}
@ -290,20 +305,13 @@ func (s *Server) Open() error {
log.Println(errors.Wrap(err, "logging startup"))
}
// Cluster settings.
s.Cluster.Broadcaster = s
s.Cluster.MaxWritesPerRequest = s.maxWritesPerRequest
// Initialize Holder.
s.holder.Broadcaster = s
// Start the BroadcastReceiver.
if err := s.BroadcastReceiver.Start(s); err != nil {
return fmt.Errorf("starting BroadcastReceiver: %v", err)
// Initialize id-key storage.
if err := s.translateFile.Open(); err != nil {
return err
}
// Open Cluster management.
if err := s.Cluster.open(); err != nil {
if err := s.Cluster.waitForStarted(); err != nil {
return fmt.Errorf("opening Cluster: %v", err)
}
@ -534,6 +542,12 @@ func (s *Server) SendTo(to *Node, pb proto.Message) error {
return s.defaultClient.SendMessage(context.Background(), &to.URI, pb)
}
// Node returns the pilosa.Node object. It is used by membership protocols to
// get this node's name(ID), location(URI), and coordinator status.
func (s *Server) Node() *Node {
return s.Cluster.Node
}
// Server implements StatusHandler.
// LocalStatus is used to periodically sync information
// between nodes. Under normal conditions, nodes should
@ -711,6 +725,11 @@ func (s *Server) monitorRuntime() {
}
}
// ReceiveEvent implements the EventHandler interface.
func (s *Server) ReceiveEvent(e *NodeEvent) error {
return s.Cluster.ReceiveEvent(e)
}
// countOpenFiles on operating systems that support lsof.
func countOpenFiles() (int, error) {
switch runtime.GOOS {

View file

@ -75,6 +75,7 @@ type Command struct {
logger loggerLogger
Handler pilosa.Handler
API *pilosa.API
ln net.Listener
serverOptions []pilosa.ServerOption
@ -246,6 +247,12 @@ func (m *Command) SetupServer() error {
primaryTranslateStore = http.NewTranslateStore(m.Config.Translation.PrimaryURL)
}
// Set Coordinator.
coordinatorOpt := pilosa.OptServerIsCoordinator(false)
if m.Config.Cluster.Coordinator || len(m.Config.Gossip.Seeds) == 0 {
coordinatorOpt = pilosa.OptServerIsCoordinator(true)
}
serverOptions := []pilosa.ServerOption{
pilosa.OptServerAntiEntropyInterval(time.Duration(m.Config.AntiEntropy.Interval)),
pilosa.OptServerLongQueryTime(time.Duration(m.Config.Cluster.LongQueryTime)),
@ -264,6 +271,7 @@ func (m *Command) SetupServer() error {
pilosa.OptServerInternalClient(http.NewInternalClientFromURI(uri, c)),
pilosa.OptServerPrimaryTranslateStore(primaryTranslateStore),
pilosa.OptServerClusterDisabled(m.Config.Cluster.Disabled, m.Config.Cluster.Hosts),
coordinatorOpt,
}
serverOptions = append(serverOptions, m.serverOptions...)
@ -274,14 +282,14 @@ func (m *Command) SetupServer() error {
return errors.Wrap(err, "new server")
}
api, err := pilosa.NewAPI(pilosa.OptAPIServer(m.Server))
m.API, err = pilosa.NewAPI(pilosa.OptAPIServer(m.Server))
if err != nil {
return errors.Wrap(err, "new api")
}
m.Handler, err = http.NewHandler(
http.OptHandlerAllowedOrigins(m.Config.Handler.AllowedOrigins),
http.OptHandlerAPI(api),
http.OptHandlerAPI(m.API),
http.OptHandlerLogger(m.logger),
http.OptHandlerListener(m.ln),
)
@ -312,19 +320,10 @@ func (m *Command) SetupNetworking() error {
}
}
// Set Coordinator.
if m.Config.Cluster.Coordinator || len(m.Config.Gossip.Seeds) == 0 {
m.Server.Cluster.Coordinator = m.Server.NodeID
m.Server.Cluster.Node.IsCoordinator = true
}
gossipEventReceiver := gossip.NewGossipEventReceiver(m.logger)
m.Server.Cluster.EventReceiver = gossipEventReceiver
gossipMemberSet, err := gossip.NewGossipMemberSet(
m.Server.NodeID,
m.Server.URI.Host(),
m.Config.Gossip,
gossipEventReceiver,
m.Server,
gossip.WithLogger(m.logger.Logger()),
gossip.WithTransport(transport),
@ -332,9 +331,7 @@ func (m *Command) SetupNetworking() error {
if err != nil {
return errors.Wrap(err, "getting memberset")
}
gossipMemberSet.Logger = m.logger
m.Server.Cluster.MemberSet = gossipMemberSet
m.Server.BroadcastReceiver = gossipMemberSet
return nil
}

File diff suppressed because one or more lines are too long

View file

@ -196,13 +196,6 @@ func (m *Main) RunWithTransport(host string, bindPort int, joinSeeds []string) (
- SetupNetworking (does the gossip or static stuff) - calls NewTransport
- Open server - calls OpenListener
*/
// SetupServer
err = m.SetupServer()
if err != nil {
return seed, err
}
// Open gossip transport to use in SetupServer.
transport, err := gossip.NewTransport(host, bindPort, nil)
if err != nil {
@ -215,21 +208,21 @@ func (m *Main) RunWithTransport(host string, bindPort int, joinSeeds []string) (
} else {
m.Config.Gossip.Seeds = []string{transport.URI.String()}
}
seed = transport.URI.String()
// SetupServer
m.Config.Cluster.Disabled = false
err = m.SetupServer()
if err != nil {
return seed, err
}
// SetupNetworking
err = m.SetupNetworking()
if err != nil {
return seed, err
}
if err = m.Server.BroadcastReceiver.Start(m.Server); err != nil {
return seed, err
}
m.Server.Cluster.Static = false
go func() {
err := m.Handler.Serve()
if err != nil {