Merge pull request #679 from linhvo/handle-action-sync

Handle actions and execute SetBit
This commit is contained in:
Linh Vo 2017-06-26 13:46:18 -05:00 • committed by GitHub
commit 8e9934a2d3
6 changed files with 408 additions and 29 deletions

View file

@ -29,7 +29,6 @@ import (
"net/http"
_ "net/http/pprof"
"os"
"sort"
"strconv"
"strings"
"time"
@ -1630,7 +1629,7 @@ func (h *Handler) handlePostInput(w http.ResponseWriter, r *http.Request) {
return
}
for _, req := range reqs {
err = h.JSONParser(req.(map[string]interface{}), index, inputDefName)
bits, err := h.JSONParser(req.(map[string]interface{}), index, inputDefName)
if err == ErrInputDefinitionNotFound {
http.Error(w, err.Error(), http.StatusNotFound)
return
@ -1638,46 +1637,67 @@ func (h *Handler) handlePostInput(w http.ResponseWriter, r *http.Request) {
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(postInputDefinitionResponse{}); err != nil {
h.logger().Printf("response encoding error: %s", err)
}
}
// JSONParser validate input json file and execute SetBit
func (h *Handler) JSONParser(req map[string]interface{}, index *Index, name string) error {
func (h *Handler) JSONParser(req map[string]interface{}, index *Index, name string) (map[string][]*Bit, error) {
inputDef := index.inputDefinition(name)
if inputDef == nil {
return ErrInputDefinitionNotFound
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 {
fmt.Errorf("field not found", key)
return nil, fmt.Errorf("field not found: %s", key)
}
}
var bits []*Bit
setBits := make(map[string][]*Bit)
for _, field := range inputDef.Fields() {
// skip field that defined in definition but not in input data
var colValue uint64
//var colValue uint64
if _, ok := req[field.Name]; !ok {
continue
} else if field.PrimaryKey {
colValue, ok := req[field.Name].(float64)
if !ok {
return fmt.Errorf("float type required, got %s:%s", field.Name, colValue)
} else {
val, ok := req[DefaultColumnLabel]
if !ok {
return errors.New("column ID not provided")
}
colValue = val.(float64)
}
}
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)
}
//bits = append(bits, bit)
setBits[frame] = append(bits, bit)
}
}
return nil
return setBits, nil
}

View file

@ -19,6 +19,11 @@ import (
"context"
"encoding/json"
"errors"
"fmt"
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/internal"
"github.com/pilosa/pilosa/pql"
"io"
"io/ioutil"
"net/http"
@ -27,11 +32,6 @@ import (
"reflect"
"strings"
"testing"
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/internal"
"github.com/pilosa/pilosa/pql"
)
// Ensure the handler returns "not found" for invalid paths.
@ -1174,3 +1174,167 @@ 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)))
fmt.Print(w.Body.String())
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)
}
}
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,26 @@ 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 {
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

@ -454,3 +454,37 @@ func TestIndex_CreateFrameWhenOpenInputDefinition(t *testing.T) {
}
}
func TestIndex_InputBits(t *testing.T) {
index := MustOpenIndex()
defer index.Close()
// Set index time quantum.
if err := index.SetTimeQuantum(pilosa.TimeQuantum("YM")); err != nil {
t.Fatal(err)
}
// Create frame.
if _, err := index.CreateFrameIfNotExists("f", pilosa.FrameOptions{}); err != nil {
t.Fatal(err)
}
var bits []*pilosa.Bit
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})
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,12 +15,13 @@
package pilosa
import (
"fmt"
"io/ioutil"
"os"
"path/filepath"
"errors"
"fmt"
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa/internal"
)
@ -267,8 +268,8 @@ type InputFrame struct {
// InputDefinitionInfo the json message format to create an InputDefinition.
type InputDefinitionInfo struct {
Frames []InputFrame `json:"frames"`
Fields []InputDefinitionField `json:"fields"`
Frames []InputFrame `json:"frames"`
Fields []InputDefinitionField `json:"fields"`
}
// Encode converts InputDefinitionInfo into its internal representation.
@ -295,7 +296,6 @@ func (i *InputDefinition) AddFrame(frame InputFrame) error {
}
return nil
}
func (i *InputDefinition) ValidateAction(action *internal.InputDefinitionAction) error {
if action.Frame == "" {
return ErrFrameRequired
@ -315,3 +315,51 @@ func (i *InputDefinition) ValidateAction(action *internal.InputDefinitionAction)
}
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 Timestams
func HandleAction(a Action, value interface{}, colID uint64) (*Bit, error) {
var err error
var bit Bit
bit.ColumnID = colID
switch a.ValueDestination {
case Mapping:
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 SingleRowBool:
switch value.(type) {
case bool:
if value.(bool) {
bit.RowID = *a.RowID
} else { // value is not True.
return nil, err
}
case float64:
if value.(float64) >= 1 {
bit.RowID = *a.RowID
} else { // value is not True.
return nil, err
}
default:
return nil, fmt.Errorf("single-row-boolean value %v must equate to a Bool", value)
}
case ValueToRow:
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) {
@ -151,3 +152,92 @@ func TestInputDefinition_LoadDefinition(t *testing.T) {
t.Fatalf("Expected frame required error, actual error: %s", err)
}
}
func TestHandleAction(t *testing.T) {
var value interface{}
colID := uint64(0)
rowID := uint64(100)
action := pilosa.Action{ValueDestination: pilosa.SingleRowBool, 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.5)
b, err = pilosa.HandleAction(action, value, colID)
if b != nil {
if b.RowID != 100 {
t.Fatalf("Unexpected rowID %v", b.RowID)
}
}
value = float64(0)
b, err = pilosa.HandleAction(action, value, colID)
if b != nil {
t.Fatalf("Expected Ignore values that do not equate to True")
}
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)
}
}
action.ValueDestination = pilosa.ValueToRow
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.Mapping
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)
}
}