Merge pull request #689 from pilosa/input-handler

Input handler
This commit is contained in:
Linh Vo 2017-06-27 17:34:48 -05:00 • committed by GitHub
commit 4d2b009315
7 changed files with 514 additions and 24 deletions

View file

@ -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
}

View file

@ -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
}

View file

@ -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)
}

View file

@ -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)
}
}

View file

@ -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
}

View file

@ -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)
}
}

View file

@ -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")