Moved more of handler to API

This commit is contained in:
Yuce Tekol 2018-03-05 18:28:24 +03:00 committed by Matthew Jaffee
parent b702f70610
commit 72526f753c
No known key found for this signature in database
GPG key ID: 51C676AF9FFCDB87
2 changed files with 141 additions and 339 deletions

2
api.go
View file

@ -484,7 +484,7 @@ func (api *API) PostClusterMessage(ctx context.Context, pb proto.Message) error
}
func (api *API) LocalID(ctx context.Context) string {
return api.Holder.LocalID
return api.Cluster.Node.ID
}
func (api *API) Schema(ctx context.Context) []*IndexInfo {

View file

@ -16,7 +16,6 @@ package pilosa
import (
"context"
"encoding/csv"
"encoding/json"
"errors"
"expvar"
@ -44,21 +43,7 @@ import (
// Handler represents an HTTP handler.
type Handler struct {
Holder *Holder
Broadcaster Broadcaster
BroadcastHandler BroadcastHandler
StatusHandler StatusHandler
FileSystem FileSystem
// Local hostname & cluster configuration.
Node *Node
Cluster *Cluster
RemoteClient *http.Client
Router *mux.Router
NormalRouter *mux.Router
RestrictedRouter *mux.Router
Router *mux.Router
// The execution engine for running queries.
Executor interface {
@ -295,8 +280,9 @@ func (h *Handler) handleWebUI(w http.ResponseWriter, r *http.Request) {
// handleGetSchema handles GET /schema requests.
func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) {
schema := h.API.Schema(r.Context())
if err := json.NewEncoder(w).Encode(getSchemaResponse{
Indexes: h.Holder.Schema(),
Indexes: schema,
}); err != nil {
h.Logger.Printf("write schema response error: %s", err)
}
@ -304,7 +290,7 @@ func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) {
// handleGetStatus handles GET /status requests.
func (h *Handler) handleGetStatus(w http.ResponseWriter, r *http.Request) {
pb, err := h.StatusHandler.ClusterStatus()
status, err := h.API.Status(r.Context())
if err != nil {
h.Logger.Printf("cluster status error: %s", err)
return
@ -774,13 +760,6 @@ func (h *Handler) handlePostFrameField(w http.ResponseWriter, r *http.Request) {
return
}
// Retrieve frame by name.
f := h.Holder.Frame(indexName, frameName)
if f == nil {
http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound)
return
}
field := &Field{
Name: fieldName,
Type: req.Type,
@ -788,23 +767,15 @@ func (h *Handler) handlePostFrameField(w http.ResponseWriter, r *http.Request) {
Max: req.Max,
}
// Create new field.
if err := f.CreateField(field); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
if err := h.API.CreateFrameField(r.Context(), indexName, frameName, field); err != nil {
if err == ErrFrameNotFound {
http.Error(w, err.Error(), http.StatusNotFound)
} else {
http.Error(w, err.Error(), http.StatusInternalServerError)
}
return
}
// Send the create field message to all nodes.
err := h.Broadcaster.SendSync(
&internal.CreateFieldMessage{
Index: indexName,
Frame: frameName,
Field: encodeField(field),
})
if err != nil {
h.Logger.Printf("problem sending CreateField message: %s", err)
}
// Encode response.
if err := json.NewEncoder(w).Encode(postFrameFieldResponse{}); err != nil {
h.Logger.Printf("response encoding error: %s", err)
@ -825,30 +796,15 @@ func (h *Handler) handleDeleteFrameField(w http.ResponseWriter, r *http.Request)
frameName := mux.Vars(r)["frame"]
fieldName := mux.Vars(r)["field"]
// Retrieve frame by name.
f := h.Holder.Frame(indexName, frameName)
if f == nil {
http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound)
if err := h.API.DeleteFrameField(r.Context(), indexName, frameName, fieldName); err != nil {
if err == ErrFrameNotFound {
http.Error(w, err.Error(), http.StatusNotFound)
} else {
http.Error(w, err.Error(), http.StatusInternalServerError)
}
return
}
// Delete field.
if err := f.DeleteField(fieldName); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Send the delete field message to all nodes.
err := h.Broadcaster.SendSync(
&internal.DeleteFieldMessage{
Index: indexName,
Frame: frameName,
Field: fieldName,
})
if err != nil {
h.Logger.Printf("problem sending DeleteField message: %s", err)
}
// Encode response.
if err := json.NewEncoder(w).Encode(deleteFrameFieldResponse{}); err != nil {
h.Logger.Printf("response encoding error: %s", err)
@ -859,24 +815,18 @@ func (h *Handler) handleGetFrameFields(w http.ResponseWriter, r *http.Request) {
indexName := mux.Vars(r)["index"]
frameName := mux.Vars(r)["frame"]
index := h.Holder.index(indexName)
if index == nil {
http.Error(w, ErrIndexNotFound.Error(), http.StatusNotFound)
return
}
frame := index.frame(frameName)
if frame == nil {
http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound)
return
}
fields, err := frame.GetFields()
if err == ErrFrameFieldsNotAllowed {
http.Error(w, err.Error(), http.StatusBadRequest)
return
} else if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
schema, err := h.API.FrameFields(r.Context(), indexName, frameName)
if err != nil {
switch err {
case ErrIndexNotFound:
fallthrough
case ErrFrameNotFound:
http.Error(w, err.Error(), http.StatusNotFound)
case ErrFrameFieldsNotAllowed:
http.Error(w, err.Error(), http.StatusBadRequest)
default:
http.Error(w, err.Error(), http.StatusInternalServerError)
}
return
}
@ -899,15 +849,16 @@ func (h *Handler) handleGetFrameViews(w http.ResponseWriter, r *http.Request) {
indexName := mux.Vars(r)["index"]
frameName := mux.Vars(r)["frame"]
// Retrieve views.
f := h.Holder.Frame(indexName, frameName)
if f == nil {
http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound)
views, err := h.API.FrameViews(r.Context(), indexName, frameName)
if err != nil {
if err == ErrFrameNotFound {
http.Error(w, err.Error(), http.StatusNotFound)
} else {
http.Error(w, err.Error(), http.StatusInternalServerError)
}
return
}
// Fetch views.
views := f.Views()
names := make([]string, len(views))
for i := range views {
names[i] = views[i].Name()
@ -925,31 +876,13 @@ func (h *Handler) handleDeleteView(w http.ResponseWriter, r *http.Request) {
frameName := mux.Vars(r)["frame"]
viewName := mux.Vars(r)["view"]
// Retrieve frame.
f := h.Holder.Frame(indexName, frameName)
if f == nil {
http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound)
return
}
// Delete the view.
if err := f.DeleteView(viewName); err != nil {
// Ingore this error because views do not exist on all nodes due to slice distribution.
if err != ErrInvalidView {
if err := h.API.DeleteView(r.Context(), indexName, frameName, viewName); err != nil {
if err == ErrFrameNotFound {
http.Error(w, err.Error(), http.StatusNotFound)
} else {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
}
// Send the delete view message to all nodes.
err := h.Broadcaster.SendSync(
&internal.DeleteViewMessage{
Index: indexName,
Frame: frameName,
View: viewName,
})
if err != nil {
h.Logger.Printf("problem sending DeleteView message: %s", err)
return
}
// Encode response.
@ -1314,28 +1247,13 @@ func (h *Handler) handleGetExportCSV(w http.ResponseWriter, r *http.Request) {
return
}
// Find the fragment.
f := h.Holder.Fragment(index, frame, view, slice)
if f == nil {
return
if err = h.API.ExportCSV(r.Context(), index, frame, view, slice, w); err != nil {
if err == ErrFragmentNotFound {
http.Error(w, err.Error(), http.StatusNotFound)
} else {
http.Error(w, err.Error(), http.StatusInternalServerError)
}
}
// Wrap writer with a CSV writer.
cw := csv.NewWriter(w)
// Iterate over each bit.
if err := f.ForEachBit(func(rowID, columnID uint64) error {
return cw.Write([]string{
strconv.FormatUint(rowID, 10),
strconv.FormatUint(columnID, 10),
})
}); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Ensure data is flushed.
cw.Flush()
}
// handleGetFragmentNodes handles /fragment/nodes requests.
@ -1351,7 +1269,7 @@ func (h *Handler) handleGetFragmentNodes(w http.ResponseWriter, r *http.Request)
}
// Retrieve fragment owner nodes.
nodes := h.Cluster.FragmentNodes(index, slice)
nodes := h.API.FragmentNodes(r.Context(), index, slice)
// Write to response.
if err := json.NewEncoder(w).Encode(nodes); err != nil {
@ -1370,9 +1288,9 @@ func (h *Handler) handleGetFragmentData(w http.ResponseWriter, r *http.Request)
}
// Retrieve fragment from holder.
f := h.Holder.Fragment(q.Get("index"), q.Get("frame"), q.Get("view"), slice)
if f == nil {
http.Error(w, "fragment not found", http.StatusNotFound)
f, err := h.API.FragmentData(r.Context(), q.Get("index"), q.Get("frame"), q.Get("view"), slice)
if err != nil {
http.Error(w, err.Error(), http.StatusNotFound)
return
}
@ -1392,31 +1310,12 @@ func (h *Handler) handlePostFragmentData(w http.ResponseWriter, r *http.Request)
return
}
// Retrieve frame.
f := h.Holder.Frame(q.Get("index"), q.Get("frame"))
if f == nil {
http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound)
return
}
// Retrieve view.
view, err := f.CreateViewIfNotExists(q.Get("view"))
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Retrieve fragment from frame.
frag, err := view.CreateFragmentIfNotExists(slice)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Read fragment in from request body.
if _, err := frag.ReadFrom(r.Body); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
if err = h.API.WriteFragmentData(r.Context(), q.Get("index"), q.Get("frame"), q.Get("view"), slice, r.Body); err != nil {
if err == ErrFrameNotFound {
http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound)
} else {
http.Error(w, err.Error(), http.StatusInternalServerError)
}
}
}
@ -1432,19 +1331,16 @@ func (h *Handler) handleGetFragmentBlockData(w http.ResponseWriter, r *http.Requ
return
}
// Retrieve fragment from holder.
f := h.Holder.Fragment(req.Index, req.Frame, req.View, req.Slice)
if f == nil {
http.Error(w, ErrFragmentNotFound.Error(), http.StatusNotFound)
resp, err := h.API.FragmentBlockData(r.Context(), req)
if err != nil {
if err == ErrFragmentNotFound {
http.Error(w, err.Error(), http.StatusNotFound)
} else {
http.Error(w, err.Error(), http.StatusInternalServerError)
}
return
}
// Read data
var resp internal.BlockDataResponse
if f != nil {
resp.RowIDs, resp.ColumnIDs = f.BlockData(int(req.Block))
}
// Encode response.
buf, err := proto.Marshal(&resp)
if err != nil {
@ -1468,16 +1364,16 @@ func (h *Handler) handleGetFragmentBlocks(w http.ResponseWriter, r *http.Request
return
}
// Retrieve fragment from holder.
f := h.Holder.Fragment(q.Get("index"), q.Get("frame"), q.Get("view"), slice)
if f == nil {
http.Error(w, "fragment not found", http.StatusNotFound)
blocks, err := h.API.FragmentBlocks(r.Context(), q.Get("index"), q.Get("frame"), q.Get("view"), slice)
if err != nil {
if err == ErrFragmentNotFound {
http.Error(w, err.Error(), http.StatusNotFound)
} else {
http.Error(w, err.Error(), http.StatusInternalServerError)
}
return
}
// Retrieve blocks.
blocks := f.Blocks()
// Encode response.
if err := json.NewEncoder(w).Encode(getFragmentBlocksResponse{
Blocks: blocks,
@ -1509,80 +1405,23 @@ func (h *Handler) handlePostFrameRestore(w http.ResponseWriter, r *http.Request)
http.Error(w, err.Error(), http.StatusBadRequest)
}
// Create a client for the remote cluster.
client := NewInternalHTTPClientFromURI(host, h.RemoteClient)
// Determine the maximum number of slices.
maxSlices, err := client.MaxSliceByIndex(r.Context())
if err != nil {
http.Error(w, "cannot determine remote slice count: "+err.Error(), http.StatusInternalServerError)
return
}
// Retrieve frame.
f := h.Holder.Frame(indexName, frameName)
if f == nil {
http.Error(w, ErrFrameNotFound.Error(), http.StatusNotFound)
return
}
// Retrieve list of all views.
views, err := client.FrameViews(r.Context(), indexName, frameName)
if err != nil {
http.Error(w, "cannot retrieve frame views: "+err.Error(), http.StatusInternalServerError)
return
}
// Loop over each slice and import it if this node owns it.
for slice := uint64(0); slice <= maxSlices[indexName]; slice++ {
// Ignore this slice if we don't own it.
if !h.Cluster.OwnsFragment(h.Node.ID, indexName, slice) {
continue
}
// Loop over view names.
for _, view := range views {
// Create view.
v, err := f.CreateViewIfNotExists(view)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Otherwise retrieve the local fragment.
frag, err := v.CreateFragmentIfNotExists(slice)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Stream backup from remote node.
rd, err := client.BackupSlice(r.Context(), indexName, frameName, view, slice)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
} else if rd == nil {
continue // slice doesn't exist
}
// Restore to local frame and always close reader.
if err := func() error {
defer rd.Close()
if _, err := frag.ReadFrom(rd); err != nil {
return err
}
return nil
}(); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
err = h.API.RestoreFrame(r.Context(), indexName, frameName, host)
switch err {
case nil:
break
case ErrFrameNotFound:
fallthrough
case ErrFragmentNotFound:
http.Error(w, err.Error(), http.StatusNotFound)
default:
http.Error(w, err.Error(), http.StatusInternalServerError)
}
}
// handleGetHosts handles /hosts requests.
func (h *Handler) handleGetHosts(w http.ResponseWriter, r *http.Request) {
if err := json.NewEncoder(w).Encode(h.Cluster.Nodes); err != nil {
hosts := h.API.ClusterHosts(r.Context())
if err := json.NewEncoder(w).Encode(hosts); err != nil {
h.Logger.Printf("write version response error: %s", err)
}
}
@ -1766,13 +1605,6 @@ func (h *Handler) handlePostInputDefinition(w http.ResponseWriter, r *http.Reque
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)
@ -1781,36 +1613,26 @@ func (h *Handler) handlePostInputDefinition(w http.ResponseWriter, r *http.Reque
return
}
if err := req.Validate(); 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 {
err = h.API.CreateInputDefinition(r.Context(), indexName, inputDefName, req)
switch err {
case nil:
break
case ErrIndexNotFound:
http.Error(w, err.Error(), http.StatusNotFound)
case ErrInputDefinitionExists:
http.Error(w, err.Error(), http.StatusConflict)
return
} else if err != nil {
case ErrInputDefinitionAttrsRequired:
fallthrough
case ErrInputDefinitionNameRequired:
fallthrough
case ErrInputDefinitionActionRequired:
fallthrough
case ErrInputDefinitionHasPrimaryKey:
fallthrough
case ErrInputDefinitionDupePrimaryKey:
http.Error(w, err.Error(), http.StatusBadRequest)
default:
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)
}
}
@ -1819,16 +1641,18 @@ func (h *Handler) handleGetInputDefinition(w http.ResponseWriter, r *http.Reques
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)
inputDef, err := h.API.InputDefinition(r.Context(), indexName, inputDefName)
if err != nil {
http.Error(w, err.Error(), http.StatusNotFound)
switch err {
case nil:
break
case ErrIndexNotFound:
fallthrough
case ErrInputDefinitionNotFound:
http.Error(w, err.Error(), http.StatusNotFound)
default:
http.Error(w, err.Error(), http.StatusInternalServerError)
}
return
}
@ -1838,7 +1662,6 @@ func (h *Handler) handleGetInputDefinition(w http.ResponseWriter, r *http.Reques
}); err != nil {
h.Logger.Printf("write status response error: %s", err)
}
}
// handleDeleteInputDefinition handles DELETE /input-definition request.
@ -1846,28 +1669,20 @@ func (h *Handler) handleDeleteInputDefinition(w http.ResponseWriter, r *http.Req
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)
if err := h.API.DeleteInputDefinition(r.Context(), indexName, inputDefName); err != nil {
switch err {
case nil:
break
case ErrIndexNotFound:
fallthrough
case ErrInputDefinitionNotFound:
http.Error(w, err.Error(), http.StatusNotFound)
default:
http.Error(w, err.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)
}
@ -1879,13 +1694,6 @@ 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)
@ -1893,28 +1701,27 @@ func (h *Handler) handlePostInput(w http.ResponseWriter, r *http.Request) {
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 {
if err = h.API.WriteInput(r.Context(), indexName, inputDefName, reqs); err != nil {
switch err {
case nil:
break
case ErrIndexNotFound:
fallthrough
case ErrInputDefinitionNotFound:
http.Error(w, err.Error(), http.StatusNotFound)
return
} else if err != nil {
default:
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
}
}
return
}
if err := json.NewEncoder(w).Encode(defaultInputDefinitionResponse{}); err != nil {
h.Logger.Printf("response encoding error: %s", err)
}
}
// <<<<<<< b702f70610962341116d9678ce7c5877b7e171ca
// handlePostClusterResizeSetCoordinator handles POST /cluster/resize/set-coordinator request.
func (h *Handler) handlePostClusterResizeSetCoordinator(w http.ResponseWriter, r *http.Request) {
// Decode request.
@ -2109,14 +1916,11 @@ func (h *Handler) InputJSONDataParser(req map[string]interface{}, index *Index,
return setBits, nil
}
// =======
// >>>>>>> Moved more of handler to API
func (h *Handler) handleRecalculateCaches(w http.ResponseWriter, r *http.Request) {
err := h.Broadcaster.SendSync(&internal.RecalculateCaches{})
if err != nil {
w.WriteHeader(http.StatusInternalServerError)
h.writeQueryResponse(w, r, &QueryResponse{Err: err})
return
}
h.Holder.RecalculateCaches()
h.API.RecalculateCaches(r.Context())
w.WriteHeader(http.StatusNoContent)
}
@ -2161,9 +1965,7 @@ func (h *Handler) handlePostClusterMessage(w http.ResponseWriter, r *http.Reques
return
}
// Forward the error message.
err = h.BroadcastHandler.ReceiveMessage(pb)
if err != nil {
if err := h.API.PostClusterMessage(r.Context(), pb); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
@ -2174,7 +1976,7 @@ func (h *Handler) handlePostClusterMessage(w http.ResponseWriter, r *http.Reques
}
func (h *Handler) handleGetID(w http.ResponseWriter, r *http.Request) {
_, err := w.Write([]byte(h.Cluster.Node.ID))
_, err := w.Write([]byte(h.API.LocalID(r.Context())))
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
}