mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-09 20:37:52 +00:00
made a container for the webapi; improved the server test
This commit is contained in:
parent
3b6e0a717f
commit
c6e1a6875d
2 changed files with 111 additions and 94 deletions
163
index/server.go
163
index/server.go
|
|
@ -12,90 +12,25 @@ import (
|
|||
"log"
|
||||
)
|
||||
|
||||
var (
|
||||
fragments[] *Fragment
|
||||
)
|
||||
type Pilosa interface{
|
||||
Union([]uint64) IBitmap
|
||||
Intersect([] uint64) IBitmap
|
||||
// SetBit(id uint64, bit_pos int64)bool
|
||||
type FragmentContainer struct {
|
||||
fragments map[string] *Fragment
|
||||
}
|
||||
func NewFragmentContainer() *FragmentContainer{
|
||||
return &FragmentContainer{make( map[string]*Fragment)}
|
||||
}
|
||||
|
||||
func (a *FragmentContainer) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
handler(w , r,a.fragments)
|
||||
}
|
||||
|
||||
type RequestJSON struct {
|
||||
Request string
|
||||
FragmentIndex int
|
||||
Args json.RawMessage
|
||||
}
|
||||
type Fragment struct {
|
||||
requestChan chan Command
|
||||
shardkey int
|
||||
impl Pilosa
|
||||
func (a *FragmentContainer) AddFragment(frame string, db string, slice int, frag_guid string) {
|
||||
f :=&Fragment{make(chan Command),frag_guid, NewGeneral(db,slice,NewMemoryStorage())}
|
||||
a.fragments[frag_guid] = f
|
||||
go f.ServeFragment()
|
||||
}
|
||||
|
||||
func (f *Fragment) ServeFragment() {
|
||||
for {
|
||||
req := <-f.requestChan
|
||||
start := time.Now()
|
||||
answer := `""`
|
||||
answer = req.Execute(f)
|
||||
delta := time.Since(start)
|
||||
var buffer bytes.Buffer
|
||||
buffer.WriteString(`{ "results":`)
|
||||
buffer.WriteString(answer)
|
||||
buffer.WriteString(fmt.Sprintf(`,"query type": "%s"`, req.QueryType()))
|
||||
buffer.WriteString(fmt.Sprintf(`, "elapsed": "%s"}`, delta))
|
||||
req.ResponseChannel() <- buffer.String()
|
||||
}
|
||||
}
|
||||
|
||||
func handler(w http.ResponseWriter, r *http.Request) {
|
||||
log.Println("GOT MESSAGE")
|
||||
if r.Method == "POST" {
|
||||
var f RequestJSON
|
||||
|
||||
bin, _ := ioutil.ReadAll(r.Body)
|
||||
err := json.Unmarshal(bin, &f)
|
||||
|
||||
if err != nil {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
fmt.Fprintf(w, fmt.Sprintf(`{ "error":"%s"}`, err))
|
||||
|
||||
}
|
||||
decoder := json.NewDecoder(bytes.NewReader(f.Args))
|
||||
request := BuildCommandFactory(&f, decoder)
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if request != nil {
|
||||
output:=`{"Error":"Invalid FragmentIndex"}`
|
||||
if f.FragmentIndex<len(fragments){
|
||||
log.Println("Sending Request")
|
||||
fc := fragments[f.FragmentIndex]
|
||||
fc.requestChan <- request
|
||||
output = request.Response()
|
||||
}
|
||||
fmt.Fprintf(w, output)
|
||||
} else {
|
||||
fmt.Fprintf(w, "NoOp")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func Startup() {
|
||||
fragments = append(fragments,&Fragment{make(chan Command),0, NewGeneral("25",0,NewMemoryStorage())})
|
||||
fragments = append(fragments,&Fragment{make(chan Command),1, NewGeneral("25",1,NewMemoryStorage())})
|
||||
fragments = append(fragments,&Fragment{make(chan Command),2, NewGeneral("25",2,NewMemoryStorage())})
|
||||
//fragments = append(fragments,&Fragment{make(chan Command),3, &Brand{}})
|
||||
for _,v:= range fragments{
|
||||
go v.ServeFragment()
|
||||
}
|
||||
}
|
||||
func Shutdown(){
|
||||
}
|
||||
|
||||
//func Add(db string, slice int, frag_type string,fragment) {
|
||||
func StartServer(port string, closeChannel chan bool,started chan bool) {
|
||||
Startup()
|
||||
fmt.Println("Ready")
|
||||
http.HandleFunc("/", handler)
|
||||
func (a *FragmentContainer) RunServer(port string, closeChannel chan bool,started chan bool) {
|
||||
http.Handle("/", a)
|
||||
|
||||
s := &http.Server{
|
||||
Addr: port,
|
||||
|
|
@ -115,8 +50,74 @@ func StartServer(port string, closeChannel chan bool,started chan bool) {
|
|||
case <-closeChannel:
|
||||
log.Printf("Server thread exit")
|
||||
l.Close()
|
||||
Shutdown()
|
||||
// Shutdown()
|
||||
return
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
type Pilosa interface{
|
||||
Union([]uint64) IBitmap
|
||||
Intersect([] uint64) IBitmap
|
||||
// SetBit(id uint64, bit_pos int64)bool
|
||||
}
|
||||
|
||||
type RequestJSON struct {
|
||||
Request string
|
||||
Fragment string
|
||||
Args json.RawMessage
|
||||
}
|
||||
type Fragment struct {
|
||||
requestChan chan Command
|
||||
FragmentGuid string
|
||||
impl Pilosa
|
||||
}
|
||||
|
||||
func (f *Fragment) ServeFragment() {
|
||||
for {
|
||||
req := <-f.requestChan
|
||||
start := time.Now()
|
||||
answer := `""`
|
||||
answer = req.Execute(f)
|
||||
delta := time.Since(start)
|
||||
var buffer bytes.Buffer
|
||||
buffer.WriteString(`{ "results":`)
|
||||
buffer.WriteString(answer)
|
||||
buffer.WriteString(fmt.Sprintf(`,"query type": "%s"`, req.QueryType()))
|
||||
buffer.WriteString(fmt.Sprintf(`, "elapsed": "%s"}`, delta))
|
||||
req.ResponseChannel() <- buffer.String()
|
||||
}
|
||||
}
|
||||
|
||||
func handler(w http.ResponseWriter, r *http.Request,fragments map[string]*Fragment) {
|
||||
log.Println("GOT MESSAGE")
|
||||
if r.Method == "POST" {
|
||||
var f RequestJSON
|
||||
|
||||
bin, _ := ioutil.ReadAll(r.Body)
|
||||
err := json.Unmarshal(bin, &f)
|
||||
|
||||
if err != nil {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
fmt.Fprintf(w, fmt.Sprintf(`{ "error":"%s"}`, err))
|
||||
|
||||
}
|
||||
decoder := json.NewDecoder(bytes.NewReader(f.Args))
|
||||
request := BuildCommandFactory(&f, decoder)
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
if request != nil {
|
||||
output:=`{"Error":"Invalid Fragment"}`
|
||||
fc,found := fragments[f.Fragment] //f.FragmentIndex<len(fragments){
|
||||
if found{
|
||||
log.Println("Sending Request")
|
||||
// fc := fragments[f.FragmentGuid]
|
||||
fc.requestChan <- request
|
||||
output = request.Response()
|
||||
}
|
||||
fmt.Fprintf(w, output)
|
||||
} else {
|
||||
fmt.Fprintf(w, "NoOp")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -5,8 +5,9 @@ import (
|
|||
. "github.com/smartystreets/goconvey/convey"
|
||||
"net/http"
|
||||
// "encoding/json"
|
||||
"io/ioutil"
|
||||
"time"
|
||||
"net/http/httptest"
|
||||
// "io/ioutil"
|
||||
// "time"
|
||||
"log"
|
||||
"strings"
|
||||
)
|
||||
|
|
@ -15,32 +16,47 @@ import (
|
|||
msg:=`
|
||||
{
|
||||
"Request": "UnionCount",
|
||||
"Fragment": 0,
|
||||
"Fragment": "AAA-BBB-CCC",
|
||||
"Args": {
|
||||
"Bitmaps":[1,2,3,4]
|
||||
}
|
||||
}`
|
||||
}`
|
||||
/*
|
||||
log.Println("POSTING:",msg)
|
||||
resp,err := http.Post("http://localhost:8089", "application/json", strings.NewReader(msg))
|
||||
defer resp.Body.Close()
|
||||
body, err := ioutil.ReadAll(resp.Body)
|
||||
log.Println(string(body),err)
|
||||
log.Println(">>>>DONE")
|
||||
*/
|
||||
|
||||
|
||||
r, err := http.NewRequest("POST", "http://api/foo", strings.NewReader(msg))
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
|
||||
w := httptest.NewRecorder()
|
||||
dummy:=FragmentContainer{make(map[string]*Fragment)}
|
||||
dummy.AddFragment("general", "25", 0, "AAA-BBB-CCC")
|
||||
|
||||
dummy.ServeHTTP(w , r )
|
||||
log.Printf("%d - %s", w.Code, w.Body.String())
|
||||
|
||||
|
||||
|
||||
}
|
||||
|
||||
func TestServer(t *testing.T) {
|
||||
Convey("Run Server", t, func() {
|
||||
Stop := make(chan bool)
|
||||
Start := make(chan bool)
|
||||
go StartServer(":8089",Stop,Start)
|
||||
select{
|
||||
case <-Start:
|
||||
// Stop := make(chan bool)
|
||||
// Start := make(chan bool)
|
||||
// go StartServer(":8089",Stop,Start)
|
||||
// select{
|
||||
// case <-Start:
|
||||
simple()
|
||||
case <-time.After(time.Duration(5) * time.Second):
|
||||
}
|
||||
Stop<- true
|
||||
// case <-time.After(time.Duration(5) * time.Second):
|
||||
// }
|
||||
// Stop<- true
|
||||
So(0, ShouldEqual, 0)
|
||||
})
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue