diff --git a/handler.go b/handler.go index bb3ed54a4..99695ee7e 100644 --- a/handler.go +++ b/handler.go @@ -115,6 +115,7 @@ func NewRouter(handler *Handler) *mux.Router { 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") @@ -1622,3 +1623,94 @@ func (h *Handler) handleDeleteInputDefinition(w http.ResponseWriter, r *http.Req } 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 validate input json file and execute SetBit +func (h *Handler) InputJsonDataParser(req map[string]interface{}, index *Index, name string) (map[string][]*Bit, error) { + inputDef := index.inputDefinition(name) + if inputDef == nil { + return nil, ErrInputDefinitionNotFound + } + // if field in input data is not in defined definition, return error + var columnLabel string + validFields := make(map[string]bool) + for _, field := range inputDef.Fields() { + validFields[field.Name] = true + if field.PrimaryKey { + columnLabel = field.Name + } + } + 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 + } + value, ok := req[columnLabel] + if !ok { + return nil, fmt.Errorf("columnLabel required") + } + colValue, ok := value.(float64) + if !ok { + return nil, fmt.Errorf("float64 require, got value:%s, type: %s", value, reflect.TypeOf(value)) + } + + for _, action := range field.Actions { + frame := action.Frame + bit, err := HandleAction(action, req[field.Name], uint64(colValue)) + 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 +} diff --git a/handler_test.go b/handler_test.go index de2831d37..05c0db4d4 100644 --- a/handler_test.go +++ b/handler_test.go @@ -1300,3 +1300,192 @@ func TestHandler_GetInputDefinition(t *testing.T) { t.Fatalf("unexpected body: %s, expect: %s", body, string(expect)) } } + +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" + + } + ] + } + ] + }` + +func TestHandler_CreateInput(t *testing.T) { + hldr := 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 + }]`) + h := NewHandler() + h.Holder = hldr.Holder + h.Cluster = NewCluster(1) + w := httptest.NewRecorder() + h.ServeHTTP(w, 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 + // f := index.Frame("cab-type") + 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 := 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"}, + } + h := NewHandler() + h.Holder = hldr.Holder + h.Cluster = NewCluster(1) + for _, test := range tests { + w := httptest.NewRecorder() + h.ServeHTTP(w, MustNewHTTPRequest("POST", "/index/i0/input/input1", bytes.NewBuffer([]byte(test.json)))) + if body := w.Body.String(); body != test.err+"\n" { + t.Fatalf("Expect error: %s, actual: %s", test.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, err := req.Encode() + def.Name = name + return def, err +} diff --git a/index.go b/index.go index ca5428e06..54a8f1433 100644 --- a/index.go +++ b/index.go @@ -755,3 +755,29 @@ func (i *Index) openInputDefinition() error { } 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(0, bit.Timestamp) + timestamps[i] = &t + } + } + + return f.Import(rowIDs, columnIDs, timestamps) +} diff --git a/index_test.go b/index_test.go index 349dc0522..aaebe288d 100644 --- a/index_test.go +++ b/index_test.go @@ -18,6 +18,7 @@ import ( "io/ioutil" "os" "reflect" + "strings" "testing" "github.com/pilosa/pilosa" @@ -459,3 +460,43 @@ func TestIndex_CreateFrameWhenOpenInputDefinition(t *testing.T) { } } + +func TestIndex_InputBits(t *testing.T) { + var bits []*pilosa.Bit + index := 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) + } +} diff --git a/input_definition.go b/input_definition.go index c83a43f5b..bef566416 100644 --- a/input_definition.go +++ b/input_definition.go @@ -15,24 +15,25 @@ package pilosa import ( + "fmt" "io/ioutil" "os" "path/filepath" "errors" - "fmt" "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa/internal" ) +// Action Mapping types const ( - Mapping = "mapping" - ValueToRow = "value-to-row" - SingleRowBool = "single-row-boolean" + InputMapping = "mapping" + InputValueToRow = "value-to-row" + InputSingleRowBool = "single-row-boolean" ) -var ValidValueDestination = []string{Mapping, ValueToRow, SingleRowBool} +var validValueDestination = []string{InputMapping, InputValueToRow, InputSingleRowBool} // InputDefinition represents a container for the data input definition. type InputDefinition struct { @@ -107,13 +108,12 @@ func (i *InputDefinition) LoadDefinition(pb *internal.InputDefinition) error { if err := i.ValidateAction(action); err != nil { return err } - if action.ValueDestination == SingleRowBool && action.Frame != "" { + if action.ValueDestination == InputSingleRowBool && action.Frame != "" { val, ok := countRowID[action.Frame] if ok && val == action.RowID { return fmt.Errorf("duplicate rowID with other field: %v", action.RowID) - } else { - countRowID[action.Frame] = action.RowID } + countRowID[action.Frame] = action.RowID } actions = append(actions, Action{ Frame: action.Frame, @@ -282,7 +282,7 @@ func (i *InputDefinitionInfo) Encode() (*internal.InputDefinition, error) { return &def, nil } -// AddFrame adds frame to input definition. +// 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 { @@ -291,23 +291,63 @@ func (i *InputDefinition) AddFrame(frame InputFrame) error { return nil } -// ValidateAction validates actions from InputDefinition. +// ValidateAction ensures the input definition action conforms to our specification. func (i *InputDefinition) ValidateAction(action *internal.InputDefinitionAction) error { if action.Frame == "" { return ErrFrameRequired } validValues := make(map[string]bool) - for _, val := range ValidValueDestination { + for _, val := range validValueDestination { validValues[val] = true } if _, ok := validValues[action.ValueDestination]; !ok { return fmt.Errorf("invalid ValueDestination: %s", action.ValueDestination) } switch action.ValueDestination { - case Mapping: + case InputMapping: if len(action.ValueMap) == 0 { return errors.New("valueMap required for map") } } 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 +// TODO handle Timestamps +func HandleAction(a Action, value interface{}, colID uint64) (*Bit, error) { + var err error + var bit Bit + bit.ColumnID = colID + + 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) + default: + return nil, fmt.Errorf("Unrecognized Value Destination: %s in Action", a.ValueDestination) + } + return &bit, err +} diff --git a/input_definition_test.go b/input_definition_test.go index 76b55c7f0..112d8eac9 100644 --- a/input_definition_test.go +++ b/input_definition_test.go @@ -18,9 +18,10 @@ import ( "encoding/json" "testing" + "strings" + "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/internal" - "strings" ) func TestInputDefinition_Open(t *testing.T) { @@ -31,8 +32,14 @@ func TestInputDefinition_Open(t *testing.T) { 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}} + 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) } @@ -113,13 +120,7 @@ func TestInputDefinition_LoadDefinition(t *testing.T) { t.Fatalf("Expected invalid ValueDestination error, actual error: %s", err) } - act := pilosa.Action{Frame: "f", ValueDestination: pilosa.SingleRowBool, ValueMap: map[string]uint64{"Green": 1}} - _, err = act.Encode() - 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) - } - - action = internal.InputDefinitionAction{Frame: "f", ValueDestination: pilosa.Mapping, RowID: 100} + action = internal.InputDefinitionAction{Frame: "f", ValueDestination: pilosa.InputMapping, RowID: 100} field = internal.InputDefinitionField{Name: "id", PrimaryKey: true, InputDefinitionActions: []*internal.InputDefinitionAction{&action}} def = &internal.InputDefinition{Name: "test", Frames: []*internal.Frame{&frames}, Fields: []*internal.InputDefinitionField{&field}} err = input.LoadDefinition(def) @@ -127,8 +128,8 @@ func TestInputDefinition_LoadDefinition(t *testing.T) { t.Fatalf("Expected valueMap required for map error, actual error: %s", err) } - action = internal.InputDefinitionAction{Frame: "f", ValueDestination: pilosa.SingleRowBool, RowID: 100} - action1 := internal.InputDefinitionAction{Frame: "f", ValueDestination: pilosa.SingleRowBool, RowID: 100} + action = internal.InputDefinitionAction{Frame: "f", ValueDestination: pilosa.InputSingleRowBool, RowID: 100} + action1 := internal.InputDefinitionAction{Frame: "f", ValueDestination: pilosa.InputSingleRowBool, RowID: 100} field1 := internal.InputDefinitionField{Name: "id", PrimaryKey: true, InputDefinitionActions: []*internal.InputDefinitionAction{&action1}} def = &internal.InputDefinition{Name: "test", Frames: []*internal.Frame{&frames}, Fields: []*internal.InputDefinitionField{&field, &field1}} err = input.LoadDefinition(def) @@ -136,10 +137,110 @@ func TestInputDefinition_LoadDefinition(t *testing.T) { t.Fatalf("Expected duplicate rowID with other field error, actual error: %s", err) } - action = internal.InputDefinitionAction{ValueDestination: pilosa.SingleRowBool, RowID: 100} + action = internal.InputDefinitionAction{ValueDestination: pilosa.InputSingleRowBool, RowID: 100} def = &internal.InputDefinition{Name: "test", Frames: []*internal.Frame{&frames}, Fields: []*internal.InputDefinitionField{&field}} err = input.LoadDefinition(def) if !strings.Contains(err.Error(), "frame required") { t.Fatalf("Expected frame required error, actual error: %s", err) } } + +func TestActionEncoding(t *testing.T) { + action := pilosa.Action{Frame: "f", ValueDestination: pilosa.InputSingleRowBool, ValueMap: map[string]uint64{"Green": 1}} + _, err := action.Encode() + 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) + } + + field := pilosa.InputDefinitionField{Name: "id", PrimaryKey: false, Actions: []pilosa.Action{action}} + info := pilosa.InputDefinitionInfo{Fields: []pilosa.InputDefinitionField{field}} + _, err = info.Encode() + 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) + } +} + +func TestHandleAction(t *testing.T) { + var value interface{} + colID := uint64(0) + rowID := uint64(100) + action := pilosa.Action{ValueDestination: pilosa.InputSingleRowBool, RowID: &rowID} + + value = 1 + b, err := pilosa.HandleAction(action, value, colID) + 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) + if b != nil { + t.Fatalf("Expected Ignore strings, only accept boolean") + } + + value = "t" + b, err = pilosa.HandleAction(action, value, colID) + 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) + 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) + if b != nil { + t.Fatalf("Expected Ignore values that do not equate to True") + } + + value = true + b, err = pilosa.HandleAction(action, value, colID) + 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) + if b != nil { + if b.RowID != 25 { + t.Fatalf("Unexpected RowID %v", b.RowID) + } + } + value = "25" + b, err = pilosa.HandleAction(action, value, colID) + 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) + if b != nil { + t.Fatalf("Expected Ignore values that are not type string") + } + + value = 25 + b, err = pilosa.HandleAction(action, value, colID) + if b != nil { + t.Fatalf("Expected Ignore values that are not type string") + } + + action.ValueDestination = "test" + b, err = pilosa.HandleAction(action, value, colID) + if !strings.Contains(err.Error(), "Unrecognized Value Destination") { + t.Fatalf("Expected Unrecognized Value Destination error, actual error: %s", err) + } + +} diff --git a/pilosa.go b/pilosa.go index 6d90c0d75..bde8cc2c7 100644 --- a/pilosa.go +++ b/pilosa.go @@ -48,6 +48,7 @@ var ( ErrInverseRangeNotAllowed = errors.New("inverse range not allowed") ErrRangeCacheNotAllowed = errors.New("range cache not allowed") ErrFrameFieldsNotAllowed = errors.New("frame fields not allowed") + ErrInputDefinitionNotFound = errors.New("input-definition not found") ErrInvalidView = errors.New("invalid view") ErrInvalidCacheType = errors.New("invalid cache type")