Merge pull request #646 from pilosa/input-definition

Input definition
This commit is contained in:
Michael Baird 2017-07-27 10:10:46 -05:00 • committed by GitHub
commit fb16f171f1
19 changed files with 3416 additions and 112 deletions

View file

@ -108,11 +108,13 @@ var NopBroadcastReceiver = &nopBroadcastReceiver{}
// Broadcast message types.
const (
MessageTypeCreateSlice = 1
MessageTypeCreateIndex = 2
MessageTypeDeleteIndex = 3
MessageTypeCreateFrame = 4
MessageTypeDeleteFrame = 5
MessageTypeCreateSlice = 1
MessageTypeCreateIndex = 2
MessageTypeDeleteIndex = 3
MessageTypeCreateFrame = 4
MessageTypeDeleteFrame = 5
MessageTypeCreateInputDefinition = 6
MessageTypeDeleteInputDefinition = 7
)
// MarshalMessage encodes the protobuf message into a byte slice.
@ -129,6 +131,10 @@ func MarshalMessage(m proto.Message) ([]byte, error) {
typ = MessageTypeCreateFrame
case *internal.DeleteFrameMessage:
typ = MessageTypeDeleteFrame
case *internal.CreateInputDefinitionMessage:
typ = MessageTypeCreateInputDefinition
case *internal.DeleteInputDefinitionMessage:
typ = MessageTypeDeleteInputDefinition
default:
return nil, fmt.Errorf("message type not implemented for marshalling: %s", reflect.TypeOf(obj))
}
@ -155,6 +161,10 @@ func UnmarshalMessage(buf []byte) (proto.Message, error) {
m = &internal.CreateFrameMessage{}
case MessageTypeDeleteFrame:
m = &internal.DeleteFrameMessage{}
case MessageTypeCreateInputDefinition:
m = &internal.CreateInputDefinitionMessage{}
case MessageTypeDeleteInputDefinition:
m = &internal.DeleteInputDefinitionMessage{}
default:
return nil, fmt.Errorf("invalid message type: %d", typ)
}

View file

@ -173,7 +173,7 @@ Response:
### Remove frame
`DELETE POST /index/<index-name>/frame/<frame-name>`
`DELETE /index/<index-name>/frame/<frame-name>`
Removes the given frame.
@ -219,6 +219,126 @@ Response:
{}
```
### Create input definition
`POST /index/<index-name>/input-definition/<input-definition-name>`
Creates an input definition in the given index with the given name.
The request payload is JSON, and it must contain the fields `frames` and `fields`. `frames` is an array of frames used within this input definition. Each frame must contain a `name` and may contain the following options:
* `rowLabel` (string): Row label of the frame.
* `timeQuantum` (string): [Time Quantum]({{< ref "data-model.md#time-quantum" >}}) for this frame.
* `inverseEnabled` (boolean): Enables [the inverted view]({{< ref "data-model.md#inverse" >}}) for this frame if `true`.
* `cacheType` (string): [ranked]({{< ref "data-model.md#ranked" >}}) or [LRU]({{< ref "data-model.md#lru" >}}) caching on this frame. Default is `lru`.
* `cacheSize` (int): Number of rows to keep in the cache. Default 50,000.
The `fields` array contains a series of JSON objects describing how to process each field received in the input data. Each `field` object must contain a `name` which maps to the source JSON field name. One field must be defined at the `primaryKey`. The `primarykey` source field name must equal the column label for the `Index`, and its value must be an unsigned integer which maps directly to a columnID in Pilosa.
* `name` (string): Maps the source data field to actions that process the field's corresponding value.
* `actions` (array): List of actions that will process the field's value.
The `action` describes how the field value will be processed. Each `action` may contain:
* `frame` (string): The Frame that will contain this action's set bits.
* `rowid` (int): The action can use this as a pre-defined SetBit rowID. The user is required to ensure this ID does not overlap with other rows in use per frame.
* `valueDestination` (string): The mapping rule used for this data.
- `value-to-row`: The value should be an integer and will map directly to a RowID.
- `single-row-boolean`: If the value is true set a bit using the `rowid`.
- `mapping`: Map the value to a RowID in the `valueMap`.
* `valueMap` (object): string and integer pairs used to map field values to RowID's.
Request:
```
curl localhost:10101/index/user/input-definition/stargazer-input \
-X POST \
-d '{
"frames":[
{
"name": "language",
"options": {"rowLabel": "language_id"}
}
],
"fields":[
{
"name": "repo_id",
"primaryKey":true
},
{
"name": "language_id",
"actions":[
{
"frame": "language",
"valueDestination": "mapping",
"valueMap": {
"Go": 5,
"Python": 17,
"C++": 10
}
}
]
}
]
}'
```
Response:
```
{}
```
### Get input definition
`GET /index/<index-name>/input-definition/<input-definition-name>`
Returns the given input definition as JSON.
Request:
```
curl -XGET localhost:10101/index/user/input-definition/stargazer-input
```
Response:
```
{"frames":[{"name":"language","options":{"rowLabel":"language_id"}}],"fields":[{"name":"repo_id","primaryKey":true},{"name":"language_id","actions":[{"frame":"language","valueDestination":"mapping","valueMap":{"Go":5,"Python":17,"C++":10}}]}]}
```
### Remove input definition
`DELETE /index/<index-name>/input-definition/<input-definition-name>`
Removes the given input definition.
Request:
```
curl -XDELETE localhost:10101/index/user/input-definition/stargazer-input
```
Response:
```
{}
```
### Process input data
`POST /index/<index-name>/input/<input-definition-name>`
Processes the JSON payload using the given input definition.
The request payload is a JSON array of objects containing one field for the primary key that corresponds to the column label, and additional fields that will be handled by corresponding actions in the input definition.
Request:
```
curl localhost:10101/index/user/input/stargazer-input \
-X POST \
-d '[{"language_id": "Go", "repo_id": 92274475}]'
```
Response:
```
{}
```
### List hosts
`GET /hosts`

View file

@ -74,9 +74,10 @@ curl localhost:10101/index/repository/frame/language \
-d '{"options": {"rowLabel": "language_id",
"inverseEnabled": true}}'
```
#### Import Some Data
The sample data for the "Star Trace" project is at [Pilosa Getting Started repository](https://github.com/pilosa/getting-started). Download the `stargazer.csv` and `language.csv` files in that repo.
#### Import Data From CSV Files
If you import data using csv files and without input defintion, download the `stargazer.csv` and `language.csv` files in that repo.
```
curl -O https://raw.githubusercontent.com/pilosa/getting-started/master/stargazer.csv
@ -100,6 +101,9 @@ docker exec -it pilosa /pilosa import -i repository -f language /language.csv
Note that, both the user IDs and the repository IDs were remapped to sequential integers in the data files, they don't correspond to actual Github IDs anymore. You can check out `language.txt` to see the mapping for languages.
### Input Definition
Alternatively Pilosa can import JSON data using an [Input Definition](../input-definition/) describing the schema and ETL rules to process the data.
#### Make Some Queries
<div class="note">

128
docs/input-definition.md Normal file
View file

@ -0,0 +1,128 @@
+++
title = "Input Definition"
+++
## Input Definition
This document builds on the data import concepts introduced in [Getting Started](../getting-started/).
Here we will demonstrate creating the index's schema and data definition. Then using this definition to import JSON data.
#### Create the Schema Using an Input Definition
Input definitions allow users to define a schema based on their data and to provide data to Pilosa in a more standard format like JSON. Once an input definition is created, we can send data to Pilosa as JSON, and as long as the data adheres to the definition, Pilosa will internally perform all of the appropriate mutations.
Before creating a schema, let's create the repository index first:
```
curl localhost:10101/index/repository \
-X POST \
-d '{"options": {"columnLabel": "repo_id"}}'
```
Then we can send the following input definition as JSON to Pilosa. The sample input defintion schema for the "Star Trace" project is at [Pilosa Getting Started repository](https://github.com/pilosa/getting-started), `input-definition.json` file
```
curl localhost:10101/index/repository/input-definition/stargazer \
-X POST \
-d '{
"frames": [
{
"name": "language",
"options": {
"inverseEnabled": true,
"timeQuantum": "YMD"
}
},
{
"name": "stargazer",
"options": {
"inverseEnabled": true,
"timeQuantum": "YMD"
}
}
],
"fields": [
{
"name": "repo_id",
"primaryKey": true
},
{
"actions": [
{
"frame": "language",
"valueDestination": "mapping",
"valueMap": {
"C": 7,
"C#": 27,
"Go": 5,
"Java": 21,
"JavaScript": 13,
"Python": 17,
}
}
],
"name": "language_id"
},
{
"actions": [
{
"frame": "stargazer",
"valueDestination": "value-to-row"
}
],
"name": "stargazer_id"
},
{
"actions": [
{
"frame": "stargazer",
"valueDestination": "set-timestamp"
}
],
"name": "time_value
}
]
}'
```
Instead of creating a `stargazer` frame and a `language` frame individually like in [Getting Started](../getting-started/), we can create multiple frames in one input definition.
We can also set `repo_id` for multiple frames at the same time by providing field actions. There are three options for valueDestination:
- value-to-row: The value for this field is used as the `rowID`.
- single-row-boolean: The value must be a boolean, and this specifies `SetBit()` or `ClearBit()`, a `rowID` must be specified for this destination type.
- mapping: The value for this field is used to lookup a `rowID` in a map. A valueMap is required for this destination type.
- set-timestamp: The value for this field is used to lookup timestamp and set timestamp for the whole frame
#### Import Data Using an Input Definition
The sample data for the "Star Trace" project is at [Pilosa Getting Started repository](https://github.com/pilosa/getting-started).
If you import data using an input definition, download the `json-input.json` file in that repo, then run the following request using the input definition created above:
```
curl localhost:10101/index/repository/input/stargazer \
-X POST \
-d '[
{
"language_id": "Go",
"repo_id": 91720568,
"stargazer_id": 513114
"time_value": "2017-05-18T20:40"
},
{
"language_id": "Python",
"repo_id": 95122322
}'
]
```
As defined in the input definition, field name `language_id` maps language to a corresponding id defined in `valueMap` and sets the appropriate bit in the `language` frame. The value corresponding to field name `stargazer_id` is added to the `stargazer` frame as rowID.
The data input above is equivalent to the following `SetBit()` operations:
```
curl localhost:10101/index/repository/query \
-X POST \
-d 'SetBit(frame="stargazer", repo_id=91720568, stargazer_id=513114)
'SetBit(frame="stargazer", repo_id=91720568, stargazer_id=513114, timestamp="2017-05-18T20:40")
SetBit(frame="language", repo_id=91720568, language_id=5)
SetBit(frame="language", repo_id=95122322, language_id=17)
'
```

View file

@ -177,6 +177,7 @@ func TestFragment_SetFieldValue(t *testing.T) {
// Limit bit depth & maximum values.
bitDepth = (bitDepth % 62) + 1
columnN = (columnN % 99) + 1
for i := range values {
values[i] = values[i] % (1 << bitDepth)
}

View file

@ -213,6 +213,11 @@ func (f *Frame) CacheSize() uint32 {
// Options returns all options for this frame.
func (f *Frame) Options() FrameOptions {
f.mu.Lock()
defer f.mu.Unlock()
return f.options()
}
func (f *Frame) options() FrameOptions {
opt := FrameOptions{
RowLabel: f.rowLabel,
InverseEnabled: f.inverseEnabled,
@ -221,7 +226,6 @@ func (f *Frame) Options() FrameOptions {
CacheSize: f.cacheSize,
TimeQuantum: f.timeQuantum,
}
f.mu.Unlock()
return opt
}
@ -329,14 +333,8 @@ func (f *Frame) loadMeta() error {
// saveMeta writes meta data for the frame.
func (f *Frame) saveMeta() error {
// Marshal metadata.
buf, err := proto.Marshal(&internal.FrameMeta{
RowLabel: f.rowLabel,
InverseEnabled: f.inverseEnabled,
RangeEnabled: f.rangeEnabled,
CacheType: f.cacheType,
CacheSize: f.cacheSize,
TimeQuantum: string(f.timeQuantum),
})
fo := f.options()
buf, err := proto.Marshal(fo.Encode())
if err != nil {
return err
}

View file

@ -95,20 +95,6 @@ func NewRouter(handler *Handler) *mux.Router {
router := mux.NewRouter()
router.HandleFunc("/", handler.handleWebUI).Methods("GET")
router.HandleFunc("/assets/{file}", handler.handleWebUI).Methods("GET")
router.HandleFunc("/index", handler.handleGetIndexes).Methods("GET")
router.HandleFunc("/index/{index}", handler.handleGetIndex).Methods("GET")
router.HandleFunc("/index/{index}", handler.handlePostIndex).Methods("POST")
router.HandleFunc("/index/{index}", handler.handleDeleteIndex).Methods("DELETE")
router.HandleFunc("/index/{index}/attr/diff", handler.handlePostIndexAttrDiff).Methods("POST")
//router.HandleFunc("/index/{index}/frame", handler.handleGetFrames).Methods("GET") // Not implemented.
router.HandleFunc("/index/{index}/frame/{frame}", handler.handlePostFrame).Methods("POST")
router.HandleFunc("/index/{index}/frame/{frame}", handler.handleDeleteFrame).Methods("DELETE")
router.HandleFunc("/index/{index}/query", handler.handlePostQuery).Methods("POST")
router.HandleFunc("/index/{index}/frame/{frame}/attr/diff", handler.handlePostFrameAttrDiff).Methods("POST")
router.HandleFunc("/index/{index}/frame/{frame}/restore", handler.handlePostFrameRestore).Methods("POST")
router.HandleFunc("/index/{index}/frame/{frame}/time-quantum", handler.handlePatchFrameTimeQuantum).Methods("PATCH")
router.HandleFunc("/index/{index}/frame/{frame}/views", handler.handleGetFrameViews).Methods("GET")
router.HandleFunc("/index/{index}/time-quantum", handler.handlePatchIndexTimeQuantum).Methods("PATCH")
router.PathPrefix("/debug/pprof/").Handler(http.DefaultServeMux).Methods("GET")
router.HandleFunc("/debug/vars", handler.handleExpvar).Methods("GET")
router.HandleFunc("/export", handler.handleGetExport).Methods("GET")
@ -118,6 +104,24 @@ func NewRouter(handler *Handler) *mux.Router {
router.HandleFunc("/fragment/data", handler.handlePostFragmentData).Methods("POST")
router.HandleFunc("/fragment/nodes", handler.handleGetFragmentNodes).Methods("GET")
router.HandleFunc("/import", handler.handlePostImport).Methods("POST")
router.HandleFunc("/index", handler.handleGetIndexes).Methods("GET")
router.HandleFunc("/index/{index}", handler.handleGetIndex).Methods("GET")
router.HandleFunc("/index/{index}", handler.handlePostIndex).Methods("POST")
router.HandleFunc("/index/{index}", handler.handleDeleteIndex).Methods("DELETE")
router.HandleFunc("/index/{index}/attr/diff", handler.handlePostIndexAttrDiff).Methods("POST")
//router.HandleFunc("/index/{index}/frame", handler.handleGetFrames).Methods("GET") // Not implemented.
router.HandleFunc("/index/{index}/frame/{frame}", handler.handlePostFrame).Methods("POST")
router.HandleFunc("/index/{index}/frame/{frame}", handler.handleDeleteFrame).Methods("DELETE")
router.HandleFunc("/index/{index}/frame/{frame}/attr/diff", handler.handlePostFrameAttrDiff).Methods("POST")
router.HandleFunc("/index/{index}/frame/{frame}/restore", handler.handlePostFrameRestore).Methods("POST")
router.HandleFunc("/index/{index}/frame/{frame}/time-quantum", handler.handlePatchFrameTimeQuantum).Methods("PATCH")
router.HandleFunc("/index/{index}/frame/{frame}/views", handler.handleGetFrameViews).Methods("GET")
router.HandleFunc("/index/{index}/input/{input-definition}", handler.handlePostInput).Methods("POST")
router.HandleFunc("/index/{index}/input-definition/{input-definition}", handler.handleGetInputDefinition).Methods("GET")
router.HandleFunc("/index/{index}/input-definition/{input-definition}", handler.handlePostInputDefinition).Methods("POST")
router.HandleFunc("/index/{index}/input-definition/{input-definition}", handler.handleDeleteInputDefinition).Methods("DELETE")
router.HandleFunc("/index/{index}/query", handler.handlePostQuery).Methods("POST")
router.HandleFunc("/index/{index}/time-quantum", handler.handlePatchIndexTimeQuantum).Methods("PATCH")
router.HandleFunc("/hosts", handler.handleGetHosts).Methods("GET")
router.HandleFunc("/schema", handler.handleGetSchema).Methods("GET")
router.HandleFunc("/slices/max", handler.handleGetSliceMax).Methods("GET")
@ -901,7 +905,7 @@ func (h *Handler) readProtobufQueryRequest(r *http.Request) (*QueryRequest, erro
func (h *Handler) readURLQueryRequest(r *http.Request) (*QueryRequest, error) {
q := r.URL.Query()
validQuery := validOptions(QueryRequest{})
for key, _ := range q {
for key := range q {
if _, ok := validQuery[key]; !ok {
return nil, errors.New("invalid query params")
}
@ -1491,3 +1495,248 @@ func errorString(err error) string {
}
return err.Error()
}
// handlePostInputDefinition handles POST /input-definition request.
func (h *Handler) handlePostInputDefinition(w http.ResponseWriter, r *http.Request) {
indexName := mux.Vars(r)["index"]
inputDefName := mux.Vars(r)["input-definition"]
// Find index.
index := h.Holder.Index(indexName)
if index == nil {
http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound)
return
}
// Decode request.
var req InputDefinitionInfo
err := json.NewDecoder(r.Body).Decode(&req)
if err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
// Validation the input definition with the curent index's ColumnLabel.
if err := req.Validate(index.ColumnLabel()); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
// Encode InputDefinition to its internal representation.
def := req.Encode()
def.Name = inputDefName
// Create InputDefinition.
_, err = index.CreateInputDefinition(def)
if err == ErrInputDefinitionExists {
http.Error(w, err.Error(), http.StatusConflict)
return
} else if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
err = h.Broadcaster.SendSync(
&internal.CreateInputDefinitionMessage{
Index: indexName,
Definition: def,
})
if err != nil {
h.logger().Printf("problem sending CreateInputDefinition message: %s", err)
}
if err := json.NewEncoder(w).Encode(defaultInputDefinitionResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
}
}
// handleGetInputDefinition handles GET /input-definition request.
func (h *Handler) handleGetInputDefinition(w http.ResponseWriter, r *http.Request) {
indexName := mux.Vars(r)["index"]
inputDefName := mux.Vars(r)["input-definition"]
// Find index.
index := h.Holder.Index(indexName)
if index == nil {
http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound)
return
}
inputDef, err := index.InputDefinition(inputDefName)
if err != nil {
http.Error(w, err.Error(), http.StatusNotFound)
return
}
if err = json.NewEncoder(w).Encode(InputDefinitionInfo{
Frames: inputDef.frames,
Fields: inputDef.fields,
}); err != nil {
h.logger().Printf("write status response error: %s", err)
}
}
// handleDeleteInputDefinition handles DELETE /input-definition request.
func (h *Handler) handleDeleteInputDefinition(w http.ResponseWriter, r *http.Request) {
indexName := mux.Vars(r)["index"]
inputDefName := mux.Vars(r)["input-definition"]
// Find index.
index := h.Holder.Index(indexName)
if index == nil {
http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound)
return
}
// Delete input definition from the index.
if err := index.DeleteInputDefinition(inputDefName); err != nil {
http.Error(w, err.Error(), http.StatusNotFound)
return
}
err := h.Broadcaster.SendSync(
&internal.DeleteInputDefinitionMessage{
Index: indexName,
Name: inputDefName,
})
if err != nil {
h.logger().Printf("problem sending DeleteInputDefinition message: %s", err)
}
if err := json.NewEncoder(w).Encode(defaultInputDefinitionResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
}
}
type defaultInputDefinitionResponse struct{}
func (h *Handler) handlePostInput(w http.ResponseWriter, r *http.Request) {
indexName := mux.Vars(r)["index"]
inputDefName := mux.Vars(r)["input-definition"]
// Find index.
index := h.Holder.Index(indexName)
if index == nil {
http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound)
return
}
// Decode request.
var reqs []interface{}
err := json.NewDecoder(r.Body).Decode(&reqs)
if err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
for _, req := range reqs {
bits, err := h.InputJSONDataParser(req.(map[string]interface{}), index, inputDefName)
if err == ErrInputDefinitionNotFound {
http.Error(w, err.Error(), http.StatusNotFound)
return
} else if err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
for fr, bs := range bits {
err := index.InputBits(fr, bs)
if err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
}
}
if err := json.NewEncoder(w).Encode(defaultInputDefinitionResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
}
}
// InputJSONDataParser validates input json file and executes SetBit.
func (h *Handler) InputJSONDataParser(req map[string]interface{}, index *Index, name string) (map[string][]*Bit, error) {
inputDef, err := index.InputDefinition(name)
if err != nil {
return nil, err
}
// If field in input data is not in defined definition, return error.
var colValue uint64
validFields := make(map[string]bool)
timestampFrame := make(map[string]int64)
for _, field := range inputDef.Fields() {
validFields[field.Name] = true
if field.PrimaryKey {
columnLabel := field.Name
value, ok := req[columnLabel]
if !ok {
return nil, fmt.Errorf("columnLabel required")
}
rawValue, ok := value.(float64) // The default JSON marshalling will interpret this as a float
if !ok {
return nil, fmt.Errorf("float64 require, got value:%s, type: %s", value, reflect.TypeOf(value))
}
colValue = uint64(rawValue)
}
// Find frame that need to add timestamp.
for _, action := range field.Actions {
if action.ValueDestination == InputSetTimestamp {
timestampFrame[action.Frame], err = GetTimeStamp(req, field.Name)
if err != nil {
return nil, err
}
}
}
}
for key := range req {
_, ok := validFields[key]
if !ok {
return nil, fmt.Errorf("field not found: %s", key)
}
}
setBits := make(map[string][]*Bit)
for _, field := range inputDef.Fields() {
// skip field that defined in definition but not in input data
if _, ok := req[field.Name]; !ok {
continue
}
// Looking into timestampFrame map and set timestamp to the whole frame
for _, action := range field.Actions {
frame := action.Frame
timestamp := timestampFrame[action.Frame]
// Skip input data field values that are set to null
if req[field.Name] == nil {
continue
}
bit, err := HandleAction(action, req[field.Name], colValue, timestamp)
if err != nil {
return nil, fmt.Errorf("error handling action: %s, err: %s", action.ValueDestination, err)
}
if bit != nil {
setBits[frame] = append(setBits[frame], bit)
}
}
}
return setBits, nil
}
// GetTimeStamp retrieves unix timestamp from Input data.
func GetTimeStamp(data map[string]interface{}, timeField string) (int64, error) {
tmstamp, ok := data[timeField]
if !ok {
return 0, nil
}
timestamp, ok := tmstamp.(string)
if !ok {
return 0, fmt.Errorf("set-timestamp value must be in time format: YYYY-MM-DD, has: %v", data[timeField])
}
v, err := time.Parse(TimeFormat, timestamp)
if err != nil {
return 0, err
}
return v.Unix(), nil
}

View file

@ -17,6 +17,7 @@ package pilosa_test
import (
"bytes"
"context"
"encoding/json"
"errors"
"io/ioutil"
"net/http"
@ -237,7 +238,7 @@ func TestHandler_Query_Args_URL(t *testing.T) {
w := httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/idx0/query?slices=0,1", strings.NewReader("Count( Bitmap( id=100))")))
if w.Code != http.StatusOK {
t.Fatalf("unexpected status code: %d", w.Code, w.Body.String())
t.Fatalf("unexpected status code: %d %s", w.Code, w.Body.String())
} else if body := w.Body.String(); body != `{"results":[100]}`+"\n" {
t.Fatalf("unexpected body: %q", body)
}
@ -961,3 +962,552 @@ func TestHandler_Expvars(t *testing.T) {
t.Fatalf("unexpected status code: %d", w.Code)
}
}
// Ensure handler can create a input definition.
func TestHandler_CreateInputDefinition(t *testing.T) {
hldr := test.MustOpenHolder()
defer hldr.Close()
hldr.MustCreateIndexIfNotExists("i0", pilosa.IndexOptions{})
inputBody := []byte(`
{
"frames":[{
"name":"event-time",
"options":{
"timeQuantum": "YMD",
"inverseEnabled": false,
"cacheType": "ranked"
}
}],
"fields": [
{
"name": "columnID",
"primaryKey": true
},
{
"name": "cabType",
"actions": [
{
"frame": "cab-type",
"valueDestination": "mapping",
"valueMap": {
"Green": 1,
"Yellow": 2
}
}
]
}
]
}`)
h := test.NewHandler()
h.Holder = hldr.Holder
h.Cluster = test.NewCluster(1)
w := httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/i0/input-definition/input1", bytes.NewBuffer(inputBody)))
if w.Code != http.StatusOK {
t.Fatalf("unexpected status code: %d", w.Code)
} else if body := w.Body.String(); body != `{}`+"\n" {
t.Fatalf("unexpected body: %s", body)
}
w = httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/i0/input-definition/input1", bytes.NewBuffer(inputBody)))
if w.Code != http.StatusConflict {
t.Fatalf("unexpected status code: %d", w.Code)
} else if body := w.Body.String(); body != pilosa.ErrInputDefinitionExists.Error()+"\n" {
t.Fatalf("unexpected body: %s", body)
}
// Test index not found.
w = httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/foo/input-definition/input2", bytes.NewBuffer(inputBody)))
if w.Code != http.StatusNotFound {
t.Fatalf("unexpected status code: %d", w.Code)
} else if body := w.Body.String(); body != pilosa.ErrIndexNotFound.Error()+"\n" {
t.Fatalf("unexpected body: %s", body)
}
}
// Ensure throwing error if there's duplicated primaryKey field.
func TestHandler_DuplicatePrimaryKey(t *testing.T) {
hldr := test.MustOpenHolder()
defer hldr.Close()
hldr.MustCreateIndexIfNotExists("i0", pilosa.IndexOptions{})
h := test.NewHandler()
h.Holder = hldr.Holder
h.Cluster = test.NewCluster(1)
//Ensure throwing error if there's duplicated primaryKey field
invalidPrimaryKey := []byte(`
{
"frames":[{
"name":"event-time",
"options":{
"timeQuantum": "YMD",
"inverseEnabled": false,
"cacheType": "ranked"
}
}],
"fields": [
{
"name": "columnID",
"primaryKey": true
},
{
"name": "columnID",
"primaryKey": true
}
]
}`)
w := httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/i0/input-definition/input2", bytes.NewBuffer(invalidPrimaryKey)))
if w.Code != http.StatusBadRequest {
t.Fatalf("unexpected status code: %d", w.Code)
} else if body := w.Body.String(); body != pilosa.ErrInputDefinitionDupePrimaryKey.Error()+"\n" {
t.Fatalf("unexpected body: %s", body)
}
// Eusure throwing error if primary field's name doesn't match columnLabel
hldr.MustCreateIndexIfNotExists("i1", pilosa.IndexOptions{ColumnLabel: "id"})
unmatchColumnBody := []byte(`
{
"frames":[{
"name":"event-time",
"options":{
"timeQuantum": "YMD",
"inverseEnabled": false,
"cacheType": "ranked"
}
}],
"fields": [
{
"name": "columnID",
"primaryKey": true
}
]
}`)
w = httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/i1/input-definition/input1", bytes.NewBuffer(unmatchColumnBody)))
if w.Code != http.StatusBadRequest {
t.Fatalf("unexpected status code: %d", w.Code)
} else if body := w.Body.String(); body != pilosa.ErrInputDefinitionColumnLabel.Error()+"\n" {
t.Fatalf("unexpected body: %s", body)
}
// Eusure throwing error if request body is invalid.
jsonErrorBody := []byte(`
{
"frames":[{
"name":"event-time",
"options":{
"timeQuantum": "YMD",
"inverseEnabled": false,
"cacheType": "ranked"
}
}],
"fields": [
{
"name": "columnID",
"primaryKey": true
}`)
w = httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/i0/input-definition/input1", bytes.NewBuffer(jsonErrorBody)))
if w.Code != http.StatusBadRequest {
t.Fatalf("unexpected status code: %d", w.Code)
} else if body := w.Body.String(); body != `unexpected EOF`+"\n" {
t.Fatalf("unexpected body: %s", body)
}
}
// Ensure handler can delete a input definition.
func TestHandler_DeleteInputDefinition(t *testing.T) {
hldr := test.MustOpenHolder()
defer hldr.Close()
h := test.NewHandler()
h.Holder = hldr.Holder
h.Cluster = test.NewCluster(1)
// Test index not found.
w := httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("DELETE", "/index/i0/input-definition/test", strings.NewReader("")))
if w.Code != http.StatusNotFound {
t.Fatalf("unexpected status code: %d", w.Code)
} else if body := w.Body.String(); body != pilosa.ErrIndexNotFound.Error()+"\n" {
t.Fatalf("unexpected body: %s", body)
}
// Test input definition is deleted.
index := hldr.MustCreateIndexIfNotExists("i0", pilosa.IndexOptions{})
frames := internal.Frame{Name: "f", Meta: &internal.FrameMeta{RowLabel: "row"}}
action := internal.InputDefinitionAction{Frame: "f", ValueDestination: "mapping", ValueMap: map[string]uint64{"Green": 1}}
fields := internal.InputDefinitionField{Name: "id", PrimaryKey: true, InputDefinitionActions: []*internal.InputDefinitionAction{&action}}
def := internal.InputDefinition{Name: "test", Frames: []*internal.Frame{&frames}, Fields: []*internal.InputDefinitionField{&fields}}
_, err := index.CreateInputDefinition(&def)
if err != nil {
t.Fatal(err)
}
// Test definition not found.
w = httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("DELETE", "/index/i0/input-definition/foo", strings.NewReader("")))
if w.Code != http.StatusNotFound {
t.Fatalf("unexpected status code: %d", w.Code)
}
w = httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("DELETE", "/index/i0/input-definition/test", strings.NewReader("")))
if w.Code != http.StatusOK {
t.Fatalf("unexpected status code: %d", w.Code)
} else if body := w.Body.String(); body != `{}`+"\n" {
t.Fatalf("unexpected body: %s", body)
}
_, err = index.InputDefinition("test")
if err != pilosa.ErrInputDefinitionNotFound {
t.Fatal(err)
}
}
// Ensure handler can get existing input definition.
func TestHandler_GetInputDefinition(t *testing.T) {
hldr := test.MustOpenHolder()
defer hldr.Close()
h := test.NewHandler()
h.Holder = hldr.Holder
h.Cluster = test.NewCluster(1)
frames := internal.Frame{Name: "f", Meta: &internal.FrameMeta{RowLabel: "row"}}
action := internal.InputDefinitionAction{Frame: "f", ValueDestination: "mapping", ValueMap: map[string]uint64{"Green": 1}}
fields := internal.InputDefinitionField{Name: "id", PrimaryKey: true, InputDefinitionActions: []*internal.InputDefinitionAction{&action}}
def := internal.InputDefinition{Name: "test", Frames: []*internal.Frame{&frames}, Fields: []*internal.InputDefinitionField{&fields}}
// Return error if index does not exist.
w := httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/index/i0/input-definition/test", strings.NewReader("")))
if w.Code != http.StatusNotFound {
t.Fatalf("unexpected status code: %d", w.Code)
} else if body := w.Body.String(); body != pilosa.ErrIndexNotFound.Error()+"\n" {
t.Fatalf("unexpected body: %s, expect: %s", body, pilosa.ErrIndexNotFound)
}
// Return existing input definition.
index := hldr.MustCreateIndexIfNotExists("i0", pilosa.IndexOptions{})
inputDef, err := index.CreateInputDefinition(&def)
if err != nil {
t.Fatal(err)
}
response := &pilosa.InputDefinitionInfo{Frames: inputDef.Frames(), Fields: inputDef.Fields()}
expect, err := json.Marshal(response)
if err != nil {
t.Fatal(err)
}
w = httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/index/i0/input-definition/test", strings.NewReader("")))
if w.Code != http.StatusOK {
t.Fatalf("unexpected status code: %d", w.Code)
} else if body := w.Body.String(); body != string(expect)+"\n" {
t.Fatalf("unexpected body: %s, expect: %s", body, string(expect))
}
// Check nonexistent definition.
w = httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/index/i0/input-definition/foo", strings.NewReader("")))
if w.Code != http.StatusNotFound {
t.Fatalf("unexpected status code: %d", w.Code)
}
}
var defaultBody = `
{
"frames":[
{
"name":"cab-type",
"options": {
"timeQuantum":"YMD",
"inverseEnabled":false,
"cacheType":"ranked"
}
},
{
"name":"add-ons",
"options": {
"timeQuantum":"YMD",
"inverseEnabled":false,
"cacheType":"ranked"
}
},
{
"name":"distance-miles",
"options": {
"timeQuantum":"YMD",
"cacheType":"ranked"
}
}
],
"fields":[
{
"name":"id",
"primaryKey":true
},
{
"name":"cabType",
"actions":[
{
"frame":"cab-type",
"valueDestination":"mapping",
"valueMap":{
"green":1,
"yellow":2
}
}
]
},
{
"name":"withPet",
"actions":[
{
"frame":"add-ons",
"valueDestination":"single-row-boolean",
"rowID":100
}
]
},
{
"name":"distanceMiles",
"actions":[
{
"frame":"distance-miles",
"valueDestination":"value-to-row"
}
]
},
{
"name":"noFrame",
"actions":[
{
"frame":"foo",
"valueDestination":"value-to-row"
}
]
},
{
"name":"null_value",
"actions":[
{
"frame":"add-ons",
"valueDestination":"value-to-row"
}
]
},
{
"name":"time_value",
"actions":[
{
"frame":"add-ons",
"valueDestination":"set-timestamp"
}
]
}
]
}`
func TestHandler_CreateInput(t *testing.T) {
hldr := test.MustOpenHolder()
defer hldr.Close()
index := hldr.MustCreateIndexIfNotExists("i0", pilosa.IndexOptions{})
defBody := []byte(defaultBody)
def, err := EncodeInputDef("input1", defBody)
if err != nil {
t.Fatal(err)
}
_, err = index.CreateInputDefinition(def)
if err != nil {
t.Fatal(err)
}
inputBody := []byte(`
[{
"id": 1,
"cabType": "yellow",
"distanceMiles": 8,
"withPet": true,
"time_value": "2017-03-20T19:35",
"null_value": null
}]`)
h := test.NewHandler()
h.Holder = hldr.Holder
h.Cluster = test.NewCluster(1)
// Return error if index does not exist.
w := httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/foo/input/input1", bytes.NewBuffer(inputBody)))
if w.Code != http.StatusNotFound {
t.Fatalf("unexpected status code: %d", w.Code)
} else if body := w.Body.String(); body != pilosa.ErrIndexNotFound.Error()+"\n" {
t.Fatalf("unexpected body: %s, expect: %s", body, pilosa.ErrIndexNotFound)
}
// Check nonexistent definition.
w = httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/i0/input/input2", bytes.NewBuffer(inputBody)))
if w.Code != http.StatusNotFound {
t.Fatalf("unexpected status code: %d", w.Code)
}
// Test successfully ingest data
w = httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/i0/input/input1", bytes.NewBuffer(inputBody)))
if w.Code != http.StatusOK {
t.Fatalf("unexpected status code: %d", w.Code)
} else if body := w.Body.String(); body != `{}`+"\n" {
t.Fatalf("unexpected body: %s", body)
}
// Verify the bits set per frame.
f0 := index.Frame("distance-miles")
v0 := f0.View(pilosa.ViewStandard)
fragment0 := v0.Fragment(0)
// Verify the distanceMiles Bit was set.
if a := fragment0.Row(8).Bits(); !reflect.DeepEqual(a, []uint64{1}) {
t.Fatalf("unexpected bits: %+v", a)
}
f1 := index.Frame("add-ons")
v1 := f1.View(pilosa.ViewStandard)
fragment1 := v1.Fragment(0)
// Verify the add-ons frame does not have a distanceMiles Bit set.
// The Input process must respect the Action Frame assignments
if a := fragment1.Row(8).Bits(); !reflect.DeepEqual(a, []uint64{}) {
t.Fatalf("unexpected bits: %+v", a)
}
// Verify the withPet Bit was set.
if a := fragment1.Row(100).Bits(); !reflect.DeepEqual(a, []uint64{1}) {
t.Fatalf("unexpected bits: %+v", a)
}
}
func TestInput_JSON(t *testing.T) {
hldr := test.MustOpenHolder()
defer hldr.Close()
index := hldr.MustCreateIndexIfNotExists("i0", pilosa.IndexOptions{})
defBody := []byte(defaultBody)
def, err := EncodeInputDef("input1", defBody)
if err != nil {
t.Fatal(err)
}
_, err = index.CreateInputDefinition(def)
if err != nil {
t.Fatal(err)
}
tests := []struct {
json string
err string
}{
{json: `[{
"id": 1,
"cabType": "yellow",
"distanceMiles": 8,
"nofield": true
}]`,
err: "field not found: nofield"},
{json: `[{
"id": "abc",
"cabType": "yellow",
"distanceMiles": 8,
"withPet": true
}]`,
err: "float64 require, got value:abc, type: string"},
{json: `[{
"cabType": "yellow",
"distanceMiles": 8,
"withPet": true
}]`,
err: "columnLabel required"},
{json: `[{
"id": 1,
"cabType": "yellow",
"distanceMiles": 8,
"withPet": true
}`,
err: "unexpected EOF"},
{json: `[{
"id": 1,
"cabType": "yellow",
"distanceMiles": 8,
"noFrame": 1
}]`,
err: "Frame not found: foo"},
{json: `[{
"id": 1,
"cabType": "yellow",
"distanceMiles": 8,
"time_value": 12345
}]`,
err: "set-timestamp value must be in time format: YYYY-MM-DD, has: 12345"},
}
h := test.NewHandler()
h.Holder = hldr.Holder
h.Cluster = test.NewCluster(1)
for _, req := range tests {
w := httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/index/i0/input/input1", bytes.NewBuffer([]byte(req.json))))
if body := w.Body.String(); body != req.err+"\n" {
t.Fatalf("Expect error: %s, actual: %s", req.err, body)
}
}
}
func EncodeInputDef(name string, body []byte) (*internal.InputDefinition, error) {
var req pilosa.InputDefinitionInfo
err := json.Unmarshal(body, &req)
if err != nil {
return nil, err
}
def := req.Encode()
def.Name = name
return def, nil
}
func TestHandler_GetTimeStamp(t *testing.T) {
data := make(map[string]interface{})
timeField := "time"
data["time"] = "2017-03-20T19:35"
val, err := pilosa.GetTimeStamp(data, timeField)
if val != 1490038500 {
t.Fatalf("Timestamp is not set correctly for %s", data["time"])
}
// Verify that an integer is not a valid time format.
data["int"] = 1490000000
val, err = pilosa.GetTimeStamp(data, "int")
if !strings.Contains(err.Error(), "set-timestamp value must be in time format") {
t.Fatalf("Expected set-timestamp value must be in time format error, actual error: %s", err)
}
// Verify reversing month and year is not valid time format.
data["time"] = "03-2017-20T19:35"
val, err = pilosa.GetTimeStamp(data, timeField)
if !strings.Contains(err.Error(), "cannot parse") {
t.Fatalf("Expected Timestamp is not set correctly, actual error: %s", err)
}
// Handle time fields that do not exist.
val, err = pilosa.GetTimeStamp(data, "test")
if val != 0 {
t.Fatalf("Expected Ignore nonexistent fields")
}
}

181
index.go
View file

@ -32,6 +32,7 @@ import (
// Default index settings.
const (
DefaultColumnLabel = "columnID"
InputDefinitionDir = ".input-definitions"
)
// Index represents a container for frames.
@ -50,13 +51,16 @@ type Index struct {
// Frames by name.
frames map[string]*Frame
// Max Slice on any node in the cluster, according to this node
// Max Slice on any node in the cluster, according to this node.
remoteMaxSlice uint64
remoteMaxInverseSlice uint64
// Column attribute storage and cache
// Column attribute storage and cache.
columnAttrStore *AttrStore
// InputDefinitions by name.
inputDefinitions map[string]*InputDefinition
broadcaster Broadcaster
Stats StatsClient
@ -71,9 +75,10 @@ func NewIndex(path, name string) (*Index, error) {
}
return &Index{
path: path,
name: name,
frames: make(map[string]*Frame),
path: path,
name: name,
frames: make(map[string]*Frame),
inputDefinitions: make(map[string]*InputDefinition),
remoteMaxSlice: 0,
remoteMaxInverseSlice: 0,
@ -150,6 +155,10 @@ func (i *Index) Open() error {
return err
}
if err := i.openInputDefinitions(); err != nil {
return err
}
return nil
}
@ -167,7 +176,7 @@ func (i *Index) openFrames() error {
}
for _, fi := range fis {
if !fi.IsDir() {
if !fi.IsDir() || fi.Name() == InputDefinitionDir {
continue
}
@ -327,6 +336,11 @@ func (i *Index) SetTimeQuantum(q TimeQuantum) error {
// FramePath returns the path to a frame in the index.
func (i *Index) FramePath(name string) string { return filepath.Join(i.path, name) }
// InputDefinitionPath returns the path to the input definition directory for the index.
func (i *Index) InputDefinitionPath() string {
return filepath.Join(i.path, InputDefinitionDir)
}
// Frame returns a frame in the index by name.
func (i *Index) Frame(name string) *Frame {
i.mu.Lock()
@ -334,8 +348,20 @@ func (i *Index) Frame(name string) *Frame {
return i.frame(name)
}
// InputDefinition returns an input definition in the index by name.
func (i *Index) InputDefinition(name string) (*InputDefinition, error) {
i.mu.Lock()
defer i.mu.Unlock()
if inputDef, ok := i.inputDefinitions[name]; ok {
return inputDef, nil
}
return nil, ErrInputDefinitionNotFound
}
func (i *Index) frame(name string) *Frame { return i.frames[name] }
func (i *Index) inputDefinition(name string) *InputDefinition { return i.inputDefinitions[name] }
// Frames returns a list of all frames in the index.
func (i *Index) Frames() []*Frame {
i.mu.Lock()
@ -587,11 +613,11 @@ type IndexOptions struct {
TimeQuantum TimeQuantum `json:"timeQuantum,omitempty"`
}
// Encode converts o into its internal representation.
func (o *IndexOptions) Encode() *internal.IndexMeta {
// Encode converts i into its internal representation.
func (i *IndexOptions) Encode() *internal.IndexMeta {
return &internal.IndexMeta{
ColumnLabel: o.ColumnLabel,
TimeQuantum: string(o.TimeQuantum),
ColumnLabel: i.ColumnLabel,
TimeQuantum: string(i.TimeQuantum),
}
}
@ -614,3 +640,138 @@ type importData struct {
RowIDs []uint64
ColumnIDs []uint64
}
// CreateInputDefinition creates a new input definition.
func (i *Index) CreateInputDefinition(pb *internal.InputDefinition) (*InputDefinition, error) {
// Ensure input definition doesn't already exist.
if i.inputDefinitions[pb.Name] != nil {
return nil, ErrInputDefinitionExists
}
return i.createInputDefinition(pb)
}
func (i *Index) createInputDefinition(pb *internal.InputDefinition) (*InputDefinition, error) {
if pb.Name == "" {
return nil, ErrInputDefinitionNameRequired
}
for _, fr := range pb.Frames {
opt := FrameOptions{
RowLabel: fr.Meta.RowLabel,
InverseEnabled: fr.Meta.InverseEnabled,
CacheType: fr.Meta.CacheType,
CacheSize: fr.Meta.CacheSize,
TimeQuantum: TimeQuantum(fr.Meta.TimeQuantum),
}
_, err := i.CreateFrame(fr.Name, opt)
if err == ErrFrameExists {
continue
} else if err != nil {
return nil, err
}
}
// Initialize input definition.
inputDef, err := i.newInputDefinition(pb.Name)
if err != nil {
return nil, err
}
if err = inputDef.LoadDefinition(pb); err != nil {
return nil, err
}
if err = inputDef.saveMeta(); err != nil {
return nil, err
}
i.inputDefinitions[pb.Name] = inputDef
return inputDef, nil
}
func (i *Index) newInputDefinition(name string) (*InputDefinition, error) {
inputDef, err := NewInputDefinition(i.InputDefinitionPath(), i.name, name)
if err != nil {
return nil, err
}
inputDef.broadcaster = i.broadcaster
return inputDef, nil
}
// DeleteInputDefinition removes an input definition from the index.
func (i *Index) DeleteInputDefinition(name string) error {
// Fail if input definition doesn't exist.
_, err := i.InputDefinition(name)
if err != nil {
return err
}
i.mu.Lock()
defer i.mu.Unlock()
// Delete input definition file.
if err := os.Remove(filepath.Join(i.InputDefinitionPath(), name)); err != nil {
return err
}
// Remove reference.
delete(i.inputDefinitions, name)
return nil
}
// openInputDefinitions opens and initializes the input definitions inside the index.
func (i *Index) openInputDefinitions() error {
inputDef, err := os.Open(i.InputDefinitionPath())
if os.IsNotExist(err) {
return nil
} else if err != nil {
return err
}
defer inputDef.Close()
inputFiles, err := inputDef.Readdir(0)
for _, file := range inputFiles {
input, err := i.newInputDefinition(file.Name())
if err != nil {
return err
}
input.Open()
i.inputDefinitions[file.Name()] = input
// Create frame if it doesn't exist.
for _, fr := range input.frames {
_, err := i.CreateFrame(fr.Name, fr.Options)
if err == ErrFrameExists {
continue
} else if err != nil {
return nil
}
}
}
return nil
}
// InputBits Process the []Bit though the Frame import process
func (i *Index) InputBits(frame string, bits []*Bit) error {
var rowIDs, columnIDs []uint64
timestamps := make([]*time.Time, len(bits))
f := i.Frame(frame)
if f == nil {
return fmt.Errorf("Frame not found: %s", frame)
}
for i, bit := range bits {
if bit == nil {
continue
}
rowIDs = append(rowIDs, bit.RowID)
columnIDs = append(columnIDs, bit.ColumnID)
// Convert timestamps to time.Time.
if bit.Timestamp > 0 {
t := time.Unix(bit.Timestamp, 0)
timestamps[i] = &t
}
}
return f.Import(rowIDs, columnIDs, timestamps)
}

View file

@ -15,10 +15,13 @@
package pilosa_test
import (
"io/ioutil"
"reflect"
"strings"
"testing"
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/internal"
"github.com/pilosa/pilosa/test"
)
@ -286,3 +289,155 @@ func TestIndex_SetTimeQuantum(t *testing.T) {
t.Fatalf("unexpected quantum (reopen): %s", q)
}
}
// Ensure index can delete a frame.
func TestIndex_InvalidName(t *testing.T) {
path, err := ioutil.TempDir("", "pilosa-index-")
if err != nil {
panic(err)
}
index, err := pilosa.NewIndex(path, "ABC")
if index != nil {
t.Fatalf("unexpected index name %v", index)
}
}
func TestIndex_CreateInputDefinition(t *testing.T) {
index := test.MustOpenIndex()
defer index.Close()
// Create Input Definition.
frames := internal.Frame{Name: "f", Meta: &internal.FrameMeta{RowLabel: "row"}}
action := internal.InputDefinitionAction{Frame: "f", ValueDestination: "mapping", ValueMap: map[string]uint64{"Green": 1}}
fields := internal.InputDefinitionField{Name: "id", PrimaryKey: true, InputDefinitionActions: []*internal.InputDefinitionAction{&action}}
def := internal.InputDefinition{Name: "test", Frames: []*internal.Frame{&frames}, Fields: []*internal.InputDefinitionField{&fields}}
inputDef, err := index.CreateInputDefinition(&def)
if err != nil {
t.Fatal(err)
} else if inputDef.Frames()[0].Name != frames.Name {
t.Fatalf("unexpected input definition frames %v", inputDef.Frames())
} else if inputDef.Fields()[0].Name != fields.Name {
t.Fatalf("unexpected input definition actions %v", inputDef.Fields())
}
}
// Ensure create input definition handle correct error
func TestIndex_CreateExistingInputDefinition(t *testing.T) {
index := test.MustOpenIndex()
defer index.Close()
//Test input definition name is required
def := internal.InputDefinition{Name: "", Frames: []*internal.Frame{}, Fields: []*internal.InputDefinitionField{}}
_, err := index.CreateInputDefinition(&def)
if err != pilosa.ErrInputDefinitionNameRequired {
t.Fatal(err)
}
// Create Input Definition.
frames := internal.Frame{Name: "f", Meta: &internal.FrameMeta{RowLabel: "row"}}
action := internal.InputDefinitionAction{Frame: "f", ValueDestination: "mapping", ValueMap: map[string]uint64{"Green": 1}}
fields := internal.InputDefinitionField{Name: "id", PrimaryKey: true, InputDefinitionActions: []*internal.InputDefinitionAction{&action}}
def = internal.InputDefinition{Name: "test", Frames: []*internal.Frame{&frames}, Fields: []*internal.InputDefinitionField{&fields}}
_, err = index.CreateInputDefinition(&def)
if err != nil {
t.Fatal(err)
}
_, err = index.CreateInputDefinition(&def)
if err != pilosa.ErrInputDefinitionExists {
t.Fatal(err)
}
}
// Ensure to delete existing input definition.
func TestIndex_DeleteInputDefinition(t *testing.T) {
index := test.MustOpenIndex()
defer index.Close()
// Create Input Definition.
frames := internal.Frame{Name: "f", Meta: &internal.FrameMeta{RowLabel: "row"}}
action := internal.InputDefinitionAction{Frame: "f", ValueDestination: "mapping", ValueMap: map[string]uint64{"Green": 1}}
fields := internal.InputDefinitionField{Name: "id", PrimaryKey: true, InputDefinitionActions: []*internal.InputDefinitionAction{&action}}
def := internal.InputDefinition{Name: "test", Frames: []*internal.Frame{&frames}, Fields: []*internal.InputDefinitionField{&fields}}
_, err := index.CreateInputDefinition(&def)
if err != nil {
t.Fatal(err)
}
_, err = index.InputDefinition("test")
if err != nil {
t.Fatal(err)
}
err = index.DeleteInputDefinition("test")
if err != nil {
t.Fatal(err)
}
_, err = index.InputDefinition("test")
if err != pilosa.ErrInputDefinitionNotFound {
t.Fatal(err)
}
}
// Ensure that frame in input definition will be created when server restart
func TestIndex_CreateFrameWhenOpenInputDefinition(t *testing.T) {
index := test.MustOpenIndex()
defer index.Close()
// Create Input Definition.
frames := internal.Frame{Name: "f", Meta: &internal.FrameMeta{RowLabel: "row"}}
action := internal.InputDefinitionAction{Frame: "f", ValueDestination: "mapping", ValueMap: map[string]uint64{"Green": 1}}
fields := internal.InputDefinitionField{Name: "id", PrimaryKey: true, InputDefinitionActions: []*internal.InputDefinitionAction{&action}}
def := internal.InputDefinition{Name: "test", Frames: []*internal.Frame{&frames}, Fields: []*internal.InputDefinitionField{&fields}}
input, err := index.CreateInputDefinition(&def)
if err != nil {
t.Fatal(err)
}
input.AddFrame(pilosa.InputFrame{Name: "f1"})
index.Reopen()
if index.Frame("f1") == nil {
t.Fatal("Frame does not created when open index")
}
}
func TestIndex_InputBits(t *testing.T) {
var bits []*pilosa.Bit
index := test.MustOpenIndex()
defer index.Close()
// Set index time quantum.
if err := index.SetTimeQuantum(pilosa.TimeQuantum("YM")); err != nil {
t.Fatal(err)
}
err := index.InputBits("f", bits)
if !strings.Contains(err.Error(), "Frame not found") {
t.Fatalf("Expected Frame not found error, actual error: %s", err)
}
// Create frame.
if _, err := index.CreateFrameIfNotExists("f", pilosa.FrameOptions{}); err != nil {
t.Fatal(err)
}
bits = append(bits, &pilosa.Bit{RowID: 0, ColumnID: 0})
bits = append(bits, &pilosa.Bit{RowID: 0, ColumnID: 1})
bits = append(bits, &pilosa.Bit{RowID: 2, ColumnID: 2, Timestamp: 1})
bits = append(bits, nil)
err = index.InputBits("f", bits)
if err != nil {
t.Fatal(err)
}
f := index.Frame("f")
v := f.View(pilosa.ViewStandard)
fragment := v.Fragment(0)
// Verify the Bits were set
if a := fragment.Row(0).Bits(); !reflect.DeepEqual(a, []uint64{0, 1}) {
t.Fatalf("unexpected bits: %+v", a)
}
}

380
input_definition.go Normal file
View file

@ -0,0 +1,380 @@
// Copyright 2017 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package pilosa
import (
"fmt"
"io/ioutil"
"os"
"path/filepath"
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa/internal"
)
// Action types.
const (
InputMapping = "mapping"
InputValueToRow = "value-to-row"
InputSingleRowBool = "single-row-boolean"
InputSetTimestamp = "set-timestamp"
)
var validValueDestination = []string{InputMapping, InputValueToRow, InputSingleRowBool, InputSetTimestamp}
// InputDefinition represents a container for the data input definition.
type InputDefinition struct {
name string
path string
index string
broadcaster Broadcaster
frames []InputFrame
fields []InputDefinitionField
}
// NewInputDefinition returns a new instance of InputDefinition.
func NewInputDefinition(path, index, name string) (*InputDefinition, error) {
err := ValidateName(name)
if err != nil {
return nil, err
}
return &InputDefinition{
path: path,
index: index,
name: name,
}, nil
}
// Frames returns frames of the input definition was initialized with.
func (i *InputDefinition) Frames() []InputFrame { return i.frames }
// Fields returns fields of the input definition was initialized with.
func (i *InputDefinition) Fields() []InputDefinitionField { return i.fields }
// Open opens and initializes the InputDefinition from file.
func (i *InputDefinition) Open() error {
if err := func() error {
if err := os.MkdirAll(i.path, 0777); err != nil {
return err
}
if err := i.loadMeta(); err != nil {
return err
}
return nil
}(); err != nil {
return err
}
return nil
}
// LoadDefinition loads the protobuf format of a definition.
func (i *InputDefinition) LoadDefinition(pb *internal.InputDefinition) error {
// Copy metadata fields.
i.name = pb.Name
for _, fr := range pb.Frames {
frameMeta := fr.Meta
inputFrame := InputFrame{
Name: fr.Name,
Options: FrameOptions{
RowLabel: frameMeta.RowLabel,
InverseEnabled: frameMeta.InverseEnabled,
CacheSize: frameMeta.CacheSize,
CacheType: frameMeta.CacheType,
TimeQuantum: TimeQuantum(frameMeta.TimeQuantum),
},
}
i.frames = append(i.frames, inputFrame)
}
for _, field := range pb.Fields {
var actions []Action
for _, action := range field.InputDefinitionActions {
actions = append(actions, Action{
Frame: action.Frame,
ValueDestination: action.ValueDestination,
ValueMap: action.ValueMap,
RowID: &action.RowID,
})
}
inputField := InputDefinitionField{
Name: field.Name,
PrimaryKey: field.PrimaryKey,
Actions: actions,
}
i.fields = append(i.fields, inputField)
}
return nil
}
func (i *InputDefinition) loadMeta() error {
var pb internal.InputDefinition
buf, err := ioutil.ReadFile(filepath.Join(i.path, i.name))
if err != nil {
return err
}
if err := proto.Unmarshal(buf, &pb); err != nil {
return err
}
return i.LoadDefinition(&pb)
}
// saveMeta writes meta data for the input definition file.
func (i *InputDefinition) saveMeta() error {
if err := os.MkdirAll(i.path, 0777); err != nil {
return err
}
var frames []*internal.Frame
for _, fr := range i.frames {
frames = append(frames, fr.Encode())
}
var fields []*internal.InputDefinitionField
for _, field := range i.fields {
fields = append(fields, field.Encode())
}
// Marshal input definition.
buf, err := proto.Marshal(&internal.InputDefinition{
Name: i.name,
Frames: frames,
Fields: fields,
})
if err != nil {
return err
}
// Write to meta file.
if err := ioutil.WriteFile(filepath.Join(i.path, i.name), buf, 0666); err != nil {
return err
}
return nil
}
// InputDefinitionField descripes a single field mapping in the InputDefinition.
type InputDefinitionField struct {
Name string `json:"name,omitempty"`
PrimaryKey bool `json:"primaryKey,omitempty"`
Actions []Action `json:"actions,omitempty"`
}
// Encode converts InputDefinitionField into its internal representation.
func (o *InputDefinitionField) Encode() *internal.InputDefinitionField {
var actions []*internal.InputDefinitionAction
for _, action := range o.Actions {
actions = append(actions, action.Encode())
}
return &internal.InputDefinitionField{
Name: o.Name,
PrimaryKey: o.PrimaryKey,
InputDefinitionActions: actions,
}
}
// Action describes the mapping method for the field in the InputDefinition.
type Action struct {
Frame string `json:"frame,omitempty"`
ValueDestination string `json:"valueDestination,omitempty"`
ValueMap map[string]uint64 `json:"valueMap,omitempty"`
RowID *uint64 `json:"rowID,omitempty"`
}
// Validate ensures the input definition action conforms to our specification.
func (a *Action) Validate() error {
if a.Frame == "" {
return ErrFrameRequired
}
if !foundItem(validValueDestination, a.ValueDestination) {
return fmt.Errorf("invalid ValueDestination: %s", a.ValueDestination)
}
switch a.ValueDestination {
case InputMapping:
if len(a.ValueMap) == 0 {
return ErrInputDefinitionValueMap
}
case InputSetTimestamp:
}
return nil
}
// Encode converts Action into its internal representation.
func (a *Action) Encode() *internal.InputDefinitionAction {
return &internal.InputDefinitionAction{
Frame: a.Frame,
ValueDestination: a.ValueDestination,
ValueMap: a.ValueMap,
RowID: convert(a.RowID),
}
}
// convert pointer to uint64.
func convert(x *uint64) uint64 {
if x != nil {
return *x
}
return 0
}
// InputFrame defines the frame used in the input definition.
type InputFrame struct {
Name string `json:"name,omitempty"`
Options FrameOptions `json:"options,omitempty"`
}
// Validate the InputFrame data.
func (i *InputFrame) Validate() error {
if err := ValidateName(i.Name); err != nil {
return err
}
// TODO frame option validation
return nil
}
// Encode converts InputFrame into its internal representation.
func (i *InputFrame) Encode() *internal.Frame {
return &internal.Frame{
Name: i.Name,
Meta: i.Options.Encode(),
}
}
// InputDefinitionInfo represents the json message format needed to create an InputDefinition.
type InputDefinitionInfo struct {
Frames []InputFrame `json:"frames"`
Fields []InputDefinitionField `json:"fields"`
}
// Validate the InputDefinitionInfo data.
func (i *InputDefinitionInfo) Validate(columnLabel string) error {
numPrimaryKey := 0
accountRowID := make(map[string]uint64)
if len(i.Frames) == 0 || len(i.Fields) == 0 {
return ErrInputDefinitionAttrsRequired
}
for _, frame := range i.Frames {
if err := frame.Validate(); err != nil {
return err
}
}
// Validate columnLabel and duplicate primaryKey.
for _, field := range i.Fields {
for _, action := range field.Actions {
if err := action.Validate(); err != nil {
return err
}
if action.ValueDestination == InputSingleRowBool {
if action.RowID == nil {
return fmt.Errorf("rowID required for single-row-boolean Field %s", field.Name)
}
val, ok := accountRowID[action.Frame]
if ok && val == convert(action.RowID) {
return fmt.Errorf("duplicate rowID with other field: %v", action.RowID)
}
accountRowID[action.Frame] = convert(action.RowID)
}
}
if field.PrimaryKey {
numPrimaryKey++
if field.Name != columnLabel {
return ErrInputDefinitionColumnLabel
}
} else if len(field.Actions) == 0 {
return ErrInputDefinitionActionRequired
}
}
if len(i.Fields) > 0 && numPrimaryKey == 0 {
return ErrInputDefinitionHasPrimaryKey
}
if numPrimaryKey > 1 {
return ErrInputDefinitionDupePrimaryKey
}
return nil
}
// Encode converts InputDefinitionInfo into its internal representation.
func (i *InputDefinitionInfo) Encode() *internal.InputDefinition {
var def internal.InputDefinition
for _, f := range i.Frames {
def.Frames = append(def.Frames, f.Encode())
}
for _, f := range i.Fields {
def.Fields = append(def.Fields, f.Encode())
}
return &def
}
// AddFrame manually add frame to input definition.
func (i *InputDefinition) AddFrame(frame InputFrame) error {
i.frames = append(i.frames, frame)
if err := i.saveMeta(); err != nil {
return err
}
return nil
}
// HandleAction Process the input data with its action and return a bit to be imported later
// Note: if the Bit should not be set then nil is returned with no error
// From the JSON marshalling the possible types are: float64, boolean, string
func HandleAction(a Action, value interface{}, colID uint64, timestamp int64) (*Bit, error) {
var err error
var bit Bit
bit.ColumnID = colID
bit.Timestamp = timestamp
switch a.ValueDestination {
case InputMapping:
v, ok := value.(string)
if !ok {
return nil, fmt.Errorf("Mapping value must be a string %v", value)
}
bit.RowID, ok = a.ValueMap[v]
if !ok {
return nil, fmt.Errorf("Value %s does not exist in definition map", v)
}
case InputSingleRowBool:
v, ok := value.(bool)
if !ok {
return nil, fmt.Errorf("single-row-boolean value %v must equate to a Bool", value)
}
if v == false { // False returns a nil error and nil bit.
return nil, err
}
bit.RowID = *a.RowID
case InputValueToRow:
v, ok := value.(float64)
if !ok {
return nil, fmt.Errorf("value-to-row value must equate to an integer %v", value)
}
bit.RowID = uint64(v)
case InputSetTimestamp:
break
default:
return nil, fmt.Errorf("Unrecognized Value Destination: %s in Action", a.ValueDestination)
}
return &bit, err
}

277
input_definition_test.go Normal file
View file

@ -0,0 +1,277 @@
// Copyright 2017 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package pilosa_test
import (
"encoding/json"
"testing"
"strings"
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/internal"
"github.com/pilosa/pilosa/test"
)
func TestInputDefinition_Open(t *testing.T) {
index := test.MustOpenIndex()
defer index.Close()
// Create Input Definition.
frames := internal.Frame{Name: "f", Meta: &internal.FrameMeta{RowLabel: "row"}}
action := internal.InputDefinitionAction{Frame: "f", ValueDestination: "mapping", ValueMap: map[string]uint64{"Green": 1}}
fields := internal.InputDefinitionField{Name: "id", PrimaryKey: true, InputDefinitionActions: []*internal.InputDefinitionAction{&action}}
def := internal.InputDefinition{Name: "^", Frames: []*internal.Frame{&frames}, Fields: []*internal.InputDefinitionField{&fields}}
inputDef, err := index.CreateInputDefinition(&def)
if !strings.Contains(err.Error(), "invalid index or frame's name") {
t.Fatalf("Expected Invalid name error, actual error: %s", err)
}
def = internal.InputDefinition{Name: "test", Frames: []*internal.Frame{&frames}, Fields: []*internal.InputDefinitionField{&fields}}
inputDef, err = index.CreateInputDefinition(&def)
if err != nil {
t.Fatal(err)
}
err = inputDef.Open()
if err != nil {
t.Fatal(err)
}
}
// Verify the InputDefinition Encoding to the internal format
func TestInputDefinition_Encoding(t *testing.T) {
inputBody := []byte(`
{
"frames":[{
"name":"event-time",
"options":{
"timeQuantum": "YMD",
"inverseEnabled": true,
"cacheType": "ranked"
}
}],
"fields": [
{
"name": "id",
"primaryKey": true
},
{
"name": "cabType",
"actions": [
{
"frame": "cab-type",
"valueDestination": "mapping",
"valueMap": {
"Green": 1,
"Yellow": 2
}
}
]
}
]
}`)
var def pilosa.InputDefinitionInfo
err := json.Unmarshal(inputBody, &def)
if err != nil {
t.Fatal(err)
}
internalDef := def.Encode()
if internalDef.Frames[0].Name != "event-time" {
t.Fatalf("unexpected frame: %v", internalDef)
} else if internalDef.Frames[0].Meta.CacheType != "ranked" {
t.Fatalf("unexpected frame meta data: %v", internalDef)
} else if len(internalDef.Fields) != 2 {
t.Fatalf("unexpected number of Fields: %d", len(internalDef.Fields))
} else if len(internalDef.Fields[1].InputDefinitionActions) != 1 {
t.Fatalf("unexpected number of Actions: %v", internalDef.Fields[1].InputDefinitionActions)
} else if internalDef.Fields[1].InputDefinitionActions[0].ValueDestination != "mapping" {
t.Fatalf("unexpected ValueDestination: %v", internalDef.Fields[1].InputDefinitionActions[0])
}
}
// Test The Action validation cases
func TestActionValidation(t *testing.T) {
rowID := uint64(100)
action := pilosa.Action{Frame: "f", ValueDestination: pilosa.InputSingleRowBool, ValueMap: map[string]uint64{"Green": 1}}
field := pilosa.InputDefinitionField{Name: "id", PrimaryKey: false, Actions: []pilosa.Action{action}}
info := pilosa.InputDefinitionInfo{Fields: []pilosa.InputDefinitionField{field}}
err := info.Validate("id")
if err != pilosa.ErrInputDefinitionAttrsRequired {
t.Fatalf("Expect error: %s, actual err: %s", pilosa.ErrInputDefinitionAttrsRequired, err)
}
frame := pilosa.InputFrame{Name: "f", Options: pilosa.FrameOptions{RowLabel: "row"}}
info = pilosa.InputDefinitionInfo{Frames: []pilosa.InputFrame{frame}, Fields: []pilosa.InputDefinitionField{field}}
err = info.Validate("id")
if !strings.Contains(err.Error(), "rowID required for single-row-boolean") {
t.Fatalf("Expected rowID required for single-row-boolean error, actual error: %s", err)
}
frame = pilosa.InputFrame{Name: "^", Options: pilosa.FrameOptions{RowLabel: "row"}}
action = pilosa.Action{Frame: "f", ValueDestination: pilosa.InputSingleRowBool, RowID: &rowID}
field = pilosa.InputDefinitionField{Name: "id", PrimaryKey: true, Actions: []pilosa.Action{action}}
info = pilosa.InputDefinitionInfo{Frames: []pilosa.InputFrame{frame}, Fields: []pilosa.InputDefinitionField{field}}
err = info.Validate("id")
if err != pilosa.ErrName {
t.Fatalf("Expect error: %s, actual err: %s", pilosa.ErrName, err)
}
frame = pilosa.InputFrame{Name: "f", Options: pilosa.FrameOptions{RowLabel: "row"}}
action = pilosa.Action{ValueDestination: pilosa.InputSingleRowBool, RowID: &rowID}
field = pilosa.InputDefinitionField{Name: "id", PrimaryKey: true, Actions: []pilosa.Action{action}}
info = pilosa.InputDefinitionInfo{Frames: []pilosa.InputFrame{frame}, Fields: []pilosa.InputDefinitionField{field}}
err = info.Validate("id")
if err != pilosa.ErrFrameRequired {
t.Fatalf("Expect error: %s, actual err: %s", pilosa.ErrFrameRequired, err)
}
action = pilosa.Action{Frame: "f", ValueDestination: pilosa.InputSingleRowBool, RowID: &rowID}
field = pilosa.InputDefinitionField{Name: "id", PrimaryKey: true, Actions: []pilosa.Action{action}}
info = pilosa.InputDefinitionInfo{Frames: []pilosa.InputFrame{frame}, Fields: []pilosa.InputDefinitionField{field}}
err = info.Validate("test")
if err != pilosa.ErrInputDefinitionColumnLabel {
t.Fatalf("Expect error: %s, actual err: %s", pilosa.ErrInputDefinitionColumnLabel, err)
}
action = pilosa.Action{Frame: "f", ValueDestination: pilosa.InputSingleRowBool, RowID: &rowID}
field = pilosa.InputDefinitionField{Name: "x", PrimaryKey: false, Actions: []pilosa.Action{action}}
info = pilosa.InputDefinitionInfo{Frames: []pilosa.InputFrame{frame}, Fields: []pilosa.InputDefinitionField{field}}
err = info.Validate("id")
if err != pilosa.ErrInputDefinitionHasPrimaryKey {
t.Fatalf("Expect error: %s, actual err: %s", pilosa.ErrInputDefinitionHasPrimaryKey, err)
}
action = pilosa.Action{Frame: "f", ValueDestination: "value-to-ROW", ValueMap: map[string]uint64{"Green": 1}}
field = pilosa.InputDefinitionField{Name: "id", PrimaryKey: true, Actions: []pilosa.Action{action}}
info = pilosa.InputDefinitionInfo{Frames: []pilosa.InputFrame{frame}, Fields: []pilosa.InputDefinitionField{field}}
err = info.Validate("id")
if !strings.Contains(err.Error(), "invalid ValueDestination") {
t.Fatalf("Expected invalid ValueDestination error, actual error: %s", err)
}
action = pilosa.Action{Frame: "f", ValueDestination: pilosa.InputMapping, RowID: &rowID}
field = pilosa.InputDefinitionField{Name: "id", PrimaryKey: true, Actions: []pilosa.Action{action}}
info = pilosa.InputDefinitionInfo{Frames: []pilosa.InputFrame{frame}, Fields: []pilosa.InputDefinitionField{field}}
err = info.Validate("id")
if err != pilosa.ErrInputDefinitionValueMap {
t.Fatalf("Expect error: %s, actual err: %s", pilosa.ErrInputDefinitionValueMap, err)
}
action = pilosa.Action{Frame: "f", ValueDestination: pilosa.InputSingleRowBool, RowID: &rowID}
field = pilosa.InputDefinitionField{Name: "test", PrimaryKey: false, Actions: []pilosa.Action{action}}
action1 := pilosa.Action{Frame: "f", ValueDestination: pilosa.InputSingleRowBool, RowID: &rowID}
field1 := pilosa.InputDefinitionField{Name: "id", PrimaryKey: true, Actions: []pilosa.Action{action1}}
info = pilosa.InputDefinitionInfo{Frames: []pilosa.InputFrame{frame}, Fields: []pilosa.InputDefinitionField{field, field1}}
err = info.Validate("id")
if !strings.Contains(err.Error(), "duplicate rowID with other field") {
t.Fatalf("Expected duplicate rowID with other field error, actual error: %s", err)
}
field = pilosa.InputDefinitionField{Name: "id", PrimaryKey: true}
field1 = pilosa.InputDefinitionField{Name: "test", PrimaryKey: false}
info = pilosa.InputDefinitionInfo{Frames: []pilosa.InputFrame{frame}, Fields: []pilosa.InputDefinitionField{field, field1}}
err = info.Validate("id")
if err != pilosa.ErrInputDefinitionActionRequired {
t.Fatalf("Expect error: %s, actual err: %s", pilosa.ErrInputDefinitionActionRequired, err)
}
}
func TestHandleAction(t *testing.T) {
var value interface{}
colID := uint64(0)
rowID := uint64(100)
action := pilosa.Action{ValueDestination: pilosa.InputSingleRowBool, RowID: &rowID}
timestamp := int64(0)
value = 1
b, err := pilosa.HandleAction(action, value, colID, timestamp)
if b != nil {
t.Fatalf("Expected integer type is not handled by single-row-boolean")
} else if !strings.Contains(err.Error(), "single-row-boolean value") {
t.Fatalf("Expected single-row-boolean value error, actual error: %s", err)
}
value = "1"
b, err = pilosa.HandleAction(action, value, colID, timestamp)
if b != nil {
t.Fatalf("Expected Ignore strings, only accept boolean")
}
value = "t"
b, err = pilosa.HandleAction(action, value, colID, timestamp)
if !strings.Contains(err.Error(), "must equate to a Bool") {
t.Fatalf("Expected Unrecognized Value Destination error, actual error: %s", err)
}
value = float64(1)
b, err = pilosa.HandleAction(action, value, colID, timestamp)
if !strings.Contains(err.Error(), "must equate to a Bool") {
t.Fatalf("Expected Unrecognized Value Destination error, actual error: %s", err)
}
value = false
b, err = pilosa.HandleAction(action, value, colID, timestamp)
if b != nil {
t.Fatalf("Expected Ignore values that do not equate to True")
}
value = true
b, err = pilosa.HandleAction(action, value, colID, timestamp)
if b != nil {
if b.ColumnID != 0 {
t.Fatalf("Unexpected ColumnID %v", b.ColumnID)
}
if b.RowID != 100 {
t.Fatalf("Unexpected rowID %v", b.RowID)
}
}
action.ValueDestination = pilosa.InputValueToRow
rowID = 101
value = float64(25.0)
b, err = pilosa.HandleAction(action, value, colID, timestamp)
if b != nil {
if b.RowID != 25 {
t.Fatalf("Unexpected RowID %v", b.RowID)
}
}
value = "25"
b, err = pilosa.HandleAction(action, value, colID, timestamp)
if b != nil {
t.Fatalf("Expected Ignore values that are not type float64")
}
action.ValueDestination = pilosa.InputMapping
value = "test"
b, err = pilosa.HandleAction(action, value, colID, timestamp)
if b != nil {
t.Fatalf("Expected Ignore values that are not type string")
}
value = 25
b, err = pilosa.HandleAction(action, value, colID, timestamp)
if b != nil {
t.Fatalf("Expected Ignore values that are not type string")
}
action.ValueDestination = "test"
b, err = pilosa.HandleAction(action, value, colID, timestamp)
if !strings.Contains(err.Error(), "Unrecognized Value Destination") {
t.Fatalf("Expected Unrecognized Value Destination error, actual error: %s", err)
}
}

File diff suppressed because it is too large Load diff

View file

@ -78,6 +78,37 @@ message Index {
uint64 MaxSlice = 3;
repeated Frame Frames = 4;
repeated uint64 Slices = 5;
repeated InputDefinition InputDefinitions = 6;
}
message InputDefinition {
string Name = 1;
repeated Frame Frames = 2;
repeated InputDefinitionField Fields = 3;
}
message InputDefinitionField {
string Name = 1;
bool PrimaryKey = 2;
repeated InputDefinitionAction InputDefinitionActions = 3;
}
message InputDefinitionAction {
string Frame = 1;
string ValueDestination = 2;
map<string, uint64> ValueMap = 3;
uint64 RowID = 4;
}
message CreateInputDefinitionMessage {
string Index = 1;
InputDefinition Definition = 3;
}
message DeleteInputDefinitionMessage {
string Index = 1;
string Name = 2;
}
message NodeStatus {

View file

@ -2536,7 +2536,7 @@ func init() { proto.RegisterFile("public.proto", fileDescriptorPublic) }
var fileDescriptorPublic = []byte{
// 563 bytes of a gzipped FileDescriptorProto
0x1f, 0x8b, 0x08, 0x00, 0x00, 0x09, 0x6e, 0x88, 0x02, 0xff, 0x8c, 0x54, 0x4b, 0x8e, 0xd3, 0x40,
0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x8c, 0x54, 0x4b, 0x8e, 0xd3, 0x40,
0x10, 0xa5, 0x63, 0xe7, 0x57, 0xf9, 0x28, 0x6a, 0xf1, 0xb1, 0x10, 0x8a, 0x2c, 0x8b, 0x85, 0x57,
0x19, 0x69, 0x38, 0x00, 0xc2, 0x49, 0x46, 0xb2, 0x10, 0x23, 0xe8, 0x0c, 0xec, 0x3d, 0x33, 0xad,
0xc1, 0x92, 0x7f, 0x74, 0xb7, 0x81, 0x1c, 0x80, 0x13, 0xb0, 0xe1, 0x06, 0x70, 0x14, 0x96, 0x1c,

View file

@ -37,6 +37,16 @@ var (
ErrFrameInverseDisabled = errors.New("frame inverse disabled")
ErrColumnRowLabelEqual = errors.New("column and row labels cannot be equal")
ErrInputDefinitionExists = errors.New("input-definition already exists")
ErrInputDefinitionHasPrimaryKey = errors.New("input-definition must contain one PrimaryKey")
ErrInputDefinitionDupePrimaryKey = errors.New("input-definition can only contain one PrimaryKey")
ErrInputDefinitionColumnLabel = errors.New("PrimaryKey field name does not match columnLabel")
ErrInputDefinitionNameRequired = errors.New("input-definition name required")
ErrInputDefinitionAttrsRequired = errors.New("frames and fields are required")
ErrInputDefinitionValueMap = errors.New("valueMap required for map")
ErrInputDefinitionActionRequired = errors.New("field definitions require an action")
ErrInputDefinitionNotFound = errors.New("input-definition not found")
ErrFieldNotFound = errors.New("field not found")
ErrFieldNameRequired = errors.New("field name required")
ErrInvalidFieldType = errors.New("invalid field type")

View file

@ -322,6 +322,18 @@ func (s *Server) ReceiveMessage(pb proto.Message) error {
if err := idx.DeleteFrame(obj.Frame); err != nil {
return err
}
case *internal.CreateInputDefinitionMessage:
idx := s.Holder.Index(obj.Index)
if idx == nil {
return fmt.Errorf("Local Index not found: %s", obj.Index)
}
idx.CreateInputDefinition(obj.Definition)
case *internal.DeleteInputDefinitionMessage:
idx := s.Holder.Index(obj.Index)
err := idx.DeleteInputDefinition(obj.Name)
if err != nil {
return err
}
}
return nil
}

View file

@ -532,6 +532,29 @@ func TestMain_SendReceiveMessage(t *testing.T) {
if maxSlices1["i"] != 2 {
t.Fatalf("unexpected maxSlice on node1: %d", maxSlices1["i"])
}
// Write input definition to the first node.
if _, err := m0.CreateDefinition("i", "test", `{
"frames": [{"name": "event-time",
"options": {
"cacheType": "ranked",
"timeQuantum": "YMD"
}}],
"fields": [{"name": "columnID",
"primaryKey": true
}]}
`); err != nil {
t.Fatal(err)
}
frame0 := m0.Server.Holder.Frame("i", "event-time")
if frame0 == nil {
t.Fatal("frame not found")
}
frame1 := m1.Server.Holder.Frame("i", "event-time")
if frame1 == nil {
t.Fatal("frame not found")
}
}
// availablePorts returns a slice of ports that can be used for testing.
@ -653,6 +676,15 @@ func (m *Main) Query(index, rawQuery, query string) (string, error) {
return resp.Body, nil
}
// CreateDefinition.
func (m *Main) CreateDefinition(index, def, query string) (string, error) {
resp := MustDo("POST", m.URL()+fmt.Sprintf("/index/%s/input-definition/%s", index, def), query)
if resp.StatusCode != http.StatusOK {
return "", fmt.Errorf("invalid status: %d, body=%s", resp.StatusCode, resp.Body)
}
return resp.Body, nil
}
// SetCommand represents a command to set a bit.
type SetCommand struct {
ID uint64

View file

@ -3,7 +3,6 @@ package test
import (
"io/ioutil"
"os"
"testing"
"github.com/pilosa/pilosa"
)
@ -77,15 +76,3 @@ func (i *Index) CreateFrameIfNotExists(name string, opt pilosa.FrameOptions) (*F
}
return &Frame{Frame: f}, nil
}
// Ensure index can delete a frame.
func TestIndex_InvalidName(t *testing.T) {
path, err := ioutil.TempDir("", "pilosa-index-")
if err != nil {
panic(err)
}
index, err := pilosa.NewIndex(path, "ABC")
if index != nil {
t.Fatalf("unexpected index name %s", index)
}
}