Merge pull request #12 from platform9/packaging-etc
Change controller and router exports to make them usable as libraries
This commit is contained in:
+3
-23
@@ -44,29 +44,9 @@ func (api *API) respondWithSuccess(w http.ResponseWriter, resp []byte) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (api *API) respondWithError(w http.ResponseWriter, err error) {
|
func (api *API) respondWithError(w http.ResponseWriter, err error) {
|
||||||
var code int
|
|
||||||
var msg string
|
|
||||||
debug.PrintStack()
|
debug.PrintStack()
|
||||||
fe, ok := err.(fission.Error)
|
code, msg := fission.GetHTTPError(err)
|
||||||
if ok {
|
log.Errorf("Error: %v: %v", code, msg)
|
||||||
msg = fe.Message
|
|
||||||
switch fe.Code {
|
|
||||||
case fission.ErrorNotFound:
|
|
||||||
code = 404
|
|
||||||
case fission.ErrorInvalidArgument:
|
|
||||||
code = 400
|
|
||||||
case fission.ErrorNoSpace:
|
|
||||||
code = 500
|
|
||||||
case fission.ErrorNotAuthorized:
|
|
||||||
code = 403
|
|
||||||
default:
|
|
||||||
code = 500
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
code = 500
|
|
||||||
msg = err.Error()
|
|
||||||
}
|
|
||||||
log.Printf("Error: %v: %v", code, msg)
|
|
||||||
http.Error(w, msg, code)
|
http.Error(w, msg, code)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -74,7 +54,7 @@ func (api *API) HomeHandler(w http.ResponseWriter, r *http.Request) {
|
|||||||
fmt.Fprintf(w, "Fission API")
|
fmt.Fprintf(w, "Fission API")
|
||||||
}
|
}
|
||||||
|
|
||||||
func (api *API) serve(port int) {
|
func (api *API) Serve(port int) {
|
||||||
r := mux.NewRouter()
|
r := mux.NewRouter()
|
||||||
r.HandleFunc("/", api.HomeHandler)
|
r.HandleFunc("/", api.HomeHandler)
|
||||||
|
|
||||||
|
|||||||
@@ -208,9 +208,9 @@ func TestMain(m *testing.M) {
|
|||||||
defer os.RemoveAll(fileStore.root)
|
defer os.RemoveAll(fileStore.root)
|
||||||
|
|
||||||
api := &API{
|
api := &API{
|
||||||
FunctionStore: FunctionStore{resourceStore: *rs},
|
FunctionStore: FunctionStore{ResourceStore: *rs},
|
||||||
HTTPTriggerStore: HTTPTriggerStore{resourceStore: *rs},
|
HTTPTriggerStore: HTTPTriggerStore{ResourceStore: *rs},
|
||||||
EnvironmentStore: EnvironmentStore{resourceStore: *rs},
|
EnvironmentStore: EnvironmentStore{ResourceStore: *rs},
|
||||||
}
|
}
|
||||||
g.client = client.New("http://localhost:8888")
|
g.client = client.New("http://localhost:8888")
|
||||||
|
|
||||||
@@ -218,7 +218,7 @@ func TestMain(m *testing.M) {
|
|||||||
ks.Delete(context.Background(), "HTTPTrigger", &etcdClient.DeleteOptions{Recursive: true})
|
ks.Delete(context.Background(), "HTTPTrigger", &etcdClient.DeleteOptions{Recursive: true})
|
||||||
ks.Delete(context.Background(), "Environment", &etcdClient.DeleteOptions{Recursive: true})
|
ks.Delete(context.Background(), "Environment", &etcdClient.DeleteOptions{Recursive: true})
|
||||||
|
|
||||||
go api.serve(8888)
|
go api.Serve(8888)
|
||||||
time.Sleep(500 * time.Millisecond)
|
time.Sleep(500 * time.Millisecond)
|
||||||
|
|
||||||
resp, err := http.Get("http://localhost:8888/")
|
resp, err := http.Get("http://localhost:8888/")
|
||||||
|
|||||||
@@ -23,17 +23,17 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type EnvironmentStore struct {
|
type EnvironmentStore struct {
|
||||||
resourceStore
|
ResourceStore
|
||||||
}
|
}
|
||||||
|
|
||||||
func (es *EnvironmentStore) Create(e *fission.Environment) (string, error) {
|
func (es *EnvironmentStore) Create(e *fission.Environment) (string, error) {
|
||||||
e.Metadata.Uid = uuid.NewV4().String()
|
e.Metadata.Uid = uuid.NewV4().String()
|
||||||
return e.Metadata.Uid, es.resourceStore.create(e)
|
return e.Metadata.Uid, es.ResourceStore.create(e)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (es *EnvironmentStore) Get(m *fission.Metadata) (*fission.Environment, error) {
|
func (es *EnvironmentStore) Get(m *fission.Metadata) (*fission.Environment, error) {
|
||||||
var e fission.Environment
|
var e fission.Environment
|
||||||
err := es.resourceStore.read(m.Name, &e)
|
err := es.ResourceStore.read(m.Name, &e)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
@@ -42,7 +42,7 @@ func (es *EnvironmentStore) Get(m *fission.Metadata) (*fission.Environment, erro
|
|||||||
|
|
||||||
func (es *EnvironmentStore) Update(e *fission.Environment) (string, error) {
|
func (es *EnvironmentStore) Update(e *fission.Environment) (string, error) {
|
||||||
e.Metadata.Uid = uuid.NewV4().String()
|
e.Metadata.Uid = uuid.NewV4().String()
|
||||||
return e.Metadata.Uid, es.resourceStore.update(e)
|
return e.Metadata.Uid, es.ResourceStore.update(e)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (es *EnvironmentStore) Delete(m fission.Metadata) error {
|
func (es *EnvironmentStore) Delete(m fission.Metadata) error {
|
||||||
@@ -50,7 +50,7 @@ func (es *EnvironmentStore) Delete(m fission.Metadata) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
return es.resourceStore.delete(typeName, m.Name)
|
return es.ResourceStore.delete(typeName, m.Name)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (es *EnvironmentStore) List() ([]fission.Environment, error) {
|
func (es *EnvironmentStore) List() ([]fission.Environment, error) {
|
||||||
@@ -59,7 +59,7 @@ func (es *EnvironmentStore) List() ([]fission.Environment, error) {
|
|||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
bufs, err := es.resourceStore.getAll(typeName)
|
bufs, err := es.ResourceStore.getAll(typeName)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -33,7 +33,7 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type (
|
type (
|
||||||
fileStore struct {
|
FileStore struct {
|
||||||
root string // abs path of root of filestore
|
root string // abs path of root of filestore
|
||||||
requestChannel chan fileStoreRequest
|
requestChannel chan fileStoreRequest
|
||||||
}
|
}
|
||||||
@@ -51,8 +51,8 @@ type (
|
|||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
func makeFileStore(path string) *fileStore {
|
func MakeFileStore(path string) *FileStore {
|
||||||
fileStore := &fileStore{
|
fileStore := &FileStore{
|
||||||
root: path,
|
root: path,
|
||||||
requestChannel: make(chan fileStoreRequest),
|
requestChannel: make(chan fileStoreRequest),
|
||||||
}
|
}
|
||||||
@@ -60,7 +60,7 @@ func makeFileStore(path string) *fileStore {
|
|||||||
return fileStore
|
return fileStore
|
||||||
}
|
}
|
||||||
|
|
||||||
func (fs *fileStore) fileStoreService() {
|
func (fs *FileStore) fileStoreService() {
|
||||||
for {
|
for {
|
||||||
req := <-fs.requestChannel
|
req := <-fs.requestChannel
|
||||||
response := &fileStoreResponse{}
|
response := &fileStoreResponse{}
|
||||||
@@ -80,7 +80,7 @@ func (fs *fileStore) fileStoreService() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (fs *fileStore) read(fileName string) ([]byte, error) {
|
func (fs *FileStore) read(fileName string) ([]byte, error) {
|
||||||
req := fileStoreRequest{
|
req := fileStoreRequest{
|
||||||
requestType: READ,
|
requestType: READ,
|
||||||
fileName: fileName,
|
fileName: fileName,
|
||||||
@@ -91,7 +91,7 @@ func (fs *fileStore) read(fileName string) ([]byte, error) {
|
|||||||
return response.fileContents, response.error
|
return response.fileContents, response.error
|
||||||
}
|
}
|
||||||
|
|
||||||
func (fs *fileStore) write(fileName string, contents []byte) error {
|
func (fs *FileStore) write(fileName string, contents []byte) error {
|
||||||
req := fileStoreRequest{
|
req := fileStoreRequest{
|
||||||
requestType: WRITE,
|
requestType: WRITE,
|
||||||
fileName: fileName,
|
fileName: fileName,
|
||||||
@@ -103,7 +103,7 @@ func (fs *fileStore) write(fileName string, contents []byte) error {
|
|||||||
return response.error
|
return response.error
|
||||||
}
|
}
|
||||||
|
|
||||||
func (fs *fileStore) delete(fileName string) error {
|
func (fs *FileStore) delete(fileName string) error {
|
||||||
req := fileStoreRequest{
|
req := fileStoreRequest{
|
||||||
requestType: DELETE,
|
requestType: DELETE,
|
||||||
fileName: fileName,
|
fileName: fileName,
|
||||||
|
|||||||
@@ -33,7 +33,7 @@ func TestFileStore(t *testing.T) {
|
|||||||
log.Printf("temp dir at %v", dir)
|
log.Printf("temp dir at %v", dir)
|
||||||
|
|
||||||
// file store
|
// file store
|
||||||
fs := makeFileStore(dir)
|
fs := MakeFileStore(dir)
|
||||||
|
|
||||||
_, err = fs.read("nonexistent")
|
_, err = fs.read("nonexistent")
|
||||||
if err == nil {
|
if err == nil {
|
||||||
|
|||||||
+16
-16
@@ -23,12 +23,12 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type FunctionStore struct {
|
type FunctionStore struct {
|
||||||
resourceStore
|
ResourceStore
|
||||||
}
|
}
|
||||||
|
|
||||||
func (fs *FunctionStore) Create(f *fission.Function) (string, error) {
|
func (fs *FunctionStore) Create(f *fission.Function) (string, error) {
|
||||||
code := []byte(f.Code)
|
code := []byte(f.Code)
|
||||||
_, uid, err := fs.resourceStore.writeFile(f.Key(), code)
|
_, uid, err := fs.ResourceStore.writeFile(f.Key(), code)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
@@ -36,9 +36,9 @@ func (fs *FunctionStore) Create(f *fission.Function) (string, error) {
|
|||||||
f.Metadata.Uid = uid
|
f.Metadata.Uid = uid
|
||||||
f.Code = ""
|
f.Code = ""
|
||||||
|
|
||||||
err = fs.resourceStore.create(f)
|
err = fs.ResourceStore.create(f)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
fs.resourceStore.deleteFile(f.Key(), uid) // ignore errors
|
fs.ResourceStore.deleteFile(f.Key(), uid) // ignore errors
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
return f.Metadata.Uid, nil
|
return f.Metadata.Uid, nil
|
||||||
@@ -46,7 +46,7 @@ func (fs *FunctionStore) Create(f *fission.Function) (string, error) {
|
|||||||
|
|
||||||
func (fs *FunctionStore) Get(m *fission.Metadata) (*fission.Function, error) {
|
func (fs *FunctionStore) Get(m *fission.Metadata) (*fission.Function, error) {
|
||||||
var f fission.Function
|
var f fission.Function
|
||||||
err := fs.resourceStore.read(m.Name, &f)
|
err := fs.ResourceStore.read(m.Name, &f)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
@@ -54,10 +54,10 @@ func (fs *FunctionStore) Get(m *fission.Metadata) (*fission.Function, error) {
|
|||||||
var code []byte
|
var code []byte
|
||||||
if len(m.Uid) > 0 {
|
if len(m.Uid) > 0 {
|
||||||
log.WithFields(log.Fields{"Uid": m.Uid}).Info("fetching by uid")
|
log.WithFields(log.Fields{"Uid": m.Uid}).Info("fetching by uid")
|
||||||
code, err = fs.resourceStore.readFile(m.Name, &m.Uid)
|
code, err = fs.ResourceStore.readFile(m.Name, &m.Uid)
|
||||||
f.Metadata = *m
|
f.Metadata = *m
|
||||||
} else {
|
} else {
|
||||||
code, err = fs.resourceStore.readFile(m.Name, nil)
|
code, err = fs.ResourceStore.readFile(m.Name, nil)
|
||||||
}
|
}
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -69,24 +69,24 @@ func (fs *FunctionStore) Get(m *fission.Metadata) (*fission.Function, error) {
|
|||||||
|
|
||||||
func (fs *FunctionStore) Update(f *fission.Function) (string, error) {
|
func (fs *FunctionStore) Update(f *fission.Function) (string, error) {
|
||||||
code := []byte(f.Code)
|
code := []byte(f.Code)
|
||||||
_, uid, err := fs.resourceStore.writeFile(f.Key(), code)
|
_, uid, err := fs.ResourceStore.writeFile(f.Key(), code)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
|
|
||||||
var fnew fission.Function
|
var fnew fission.Function
|
||||||
err = fs.resourceStore.read(f.Metadata.Name, &fnew)
|
err = fs.ResourceStore.read(f.Metadata.Name, &fnew)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
fs.resourceStore.deleteFile(f.Key(), uid) // ignore err
|
fs.ResourceStore.deleteFile(f.Key(), uid) // ignore err
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
|
|
||||||
fnew.Metadata.Uid = uid
|
fnew.Metadata.Uid = uid
|
||||||
fnew.Environment = f.Environment
|
fnew.Environment = f.Environment
|
||||||
|
|
||||||
err = fs.resourceStore.update(fnew)
|
err = fs.ResourceStore.update(fnew)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
fs.resourceStore.deleteFile(f.Key(), uid) // ignore err
|
fs.ResourceStore.deleteFile(f.Key(), uid) // ignore err
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
return uid, err
|
return uid, err
|
||||||
@@ -94,12 +94,12 @@ func (fs *FunctionStore) Update(f *fission.Function) (string, error) {
|
|||||||
|
|
||||||
func (fs *FunctionStore) Delete(m fission.Metadata) error {
|
func (fs *FunctionStore) Delete(m fission.Metadata) error {
|
||||||
if len(m.Uid) == 0 {
|
if len(m.Uid) == 0 {
|
||||||
err := fs.resourceStore.deleteAllFiles(m.Name)
|
err := fs.ResourceStore.deleteAllFiles(m.Name)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
err := fs.resourceStore.deleteFile(m.Name, m.Uid)
|
err := fs.ResourceStore.deleteFile(m.Name, m.Uid)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -108,7 +108,7 @@ func (fs *FunctionStore) Delete(m fission.Metadata) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
return fs.resourceStore.delete(typeName, m.Name)
|
return fs.ResourceStore.delete(typeName, m.Name)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (fs *FunctionStore) List() ([]fission.Function, error) {
|
func (fs *FunctionStore) List() ([]fission.Function, error) {
|
||||||
@@ -117,7 +117,7 @@ func (fs *FunctionStore) List() ([]fission.Function, error) {
|
|||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
bufs, err := fs.resourceStore.getAll(typeName)
|
bufs, err := fs.ResourceStore.getAll(typeName)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -23,17 +23,17 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type HTTPTriggerStore struct {
|
type HTTPTriggerStore struct {
|
||||||
resourceStore
|
ResourceStore
|
||||||
}
|
}
|
||||||
|
|
||||||
func (hts *HTTPTriggerStore) Create(ht *fission.HTTPTrigger) (string, error) {
|
func (hts *HTTPTriggerStore) Create(ht *fission.HTTPTrigger) (string, error) {
|
||||||
ht.Metadata.Uid = uuid.NewV4().String()
|
ht.Metadata.Uid = uuid.NewV4().String()
|
||||||
return ht.Metadata.Uid, hts.resourceStore.create(ht)
|
return ht.Metadata.Uid, hts.ResourceStore.create(ht)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (hts *HTTPTriggerStore) Get(m *fission.Metadata) (*fission.HTTPTrigger, error) {
|
func (hts *HTTPTriggerStore) Get(m *fission.Metadata) (*fission.HTTPTrigger, error) {
|
||||||
var ht fission.HTTPTrigger
|
var ht fission.HTTPTrigger
|
||||||
err := hts.resourceStore.read(m.Name, &ht)
|
err := hts.ResourceStore.read(m.Name, &ht)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
@@ -42,7 +42,7 @@ func (hts *HTTPTriggerStore) Get(m *fission.Metadata) (*fission.HTTPTrigger, err
|
|||||||
|
|
||||||
func (hts *HTTPTriggerStore) Update(ht *fission.HTTPTrigger) (string, error) {
|
func (hts *HTTPTriggerStore) Update(ht *fission.HTTPTrigger) (string, error) {
|
||||||
ht.Metadata.Uid = uuid.NewV4().String()
|
ht.Metadata.Uid = uuid.NewV4().String()
|
||||||
return ht.Metadata.Uid, hts.resourceStore.update(ht)
|
return ht.Metadata.Uid, hts.ResourceStore.update(ht)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (hts *HTTPTriggerStore) Delete(m fission.Metadata) error {
|
func (hts *HTTPTriggerStore) Delete(m fission.Metadata) error {
|
||||||
@@ -50,7 +50,7 @@ func (hts *HTTPTriggerStore) Delete(m fission.Metadata) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
return hts.resourceStore.delete(typeName, m.Name)
|
return hts.ResourceStore.delete(typeName, m.Name)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (hts *HTTPTriggerStore) List() ([]fission.HTTPTrigger, error) {
|
func (hts *HTTPTriggerStore) List() ([]fission.HTTPTrigger, error) {
|
||||||
@@ -59,7 +59,7 @@ func (hts *HTTPTriggerStore) List() ([]fission.HTTPTrigger, error) {
|
|||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
bufs, err := hts.resourceStore.getAll(typeName)
|
bufs, err := hts.ResourceStore.getAll(typeName)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|||||||
+27
-21
@@ -28,18 +28,23 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type (
|
type (
|
||||||
resourceStore struct {
|
ResourceStore struct {
|
||||||
*fileStore
|
*FileStore
|
||||||
client.KeysAPI
|
client.KeysAPI
|
||||||
serializer
|
serializer
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
func makeResourceStore(fs *fileStore, ks client.KeysAPI, s serializer) *resourceStore {
|
func MakeResourceStore(fs *FileStore, etcdUrls []string) (*ResourceStore, error) {
|
||||||
return &resourceStore{fileStore: fs, KeysAPI: ks, serializer: s}
|
ks, err := getEtcdKeyAPI(etcdUrls)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
s := JsonSerializer{}
|
||||||
|
return &ResourceStore{FileStore: fs, KeysAPI: ks, serializer: s}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func getEtcdKeyAPI(etcdUrls []string) client.KeysAPI {
|
func getEtcdKeyAPI(etcdUrls []string) (client.KeysAPI, error) {
|
||||||
cfg := client.Config{
|
cfg := client.Config{
|
||||||
Endpoints: etcdUrls,
|
Endpoints: etcdUrls,
|
||||||
Transport: client.DefaultTransport,
|
Transport: client.DefaultTransport,
|
||||||
@@ -48,9 +53,10 @@ func getEtcdKeyAPI(etcdUrls []string) client.KeysAPI {
|
|||||||
}
|
}
|
||||||
c, err := client.New(cfg)
|
c, err := client.New(cfg)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Fatalf("failed to connect to etcd: %v", err)
|
log.Printf("failed to connect to etcd: %v", err)
|
||||||
|
return nil, err
|
||||||
}
|
}
|
||||||
return client.NewKeysAPI(c)
|
return client.NewKeysAPI(c), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func getTypeName(r resource) (string, error) {
|
func getTypeName(r resource) (string, error) {
|
||||||
@@ -74,7 +80,7 @@ func getKey(r resource) (string, error) {
|
|||||||
return (typName + "/" + rkey), nil
|
return (typName + "/" + rkey), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (rs *resourceStore) create(r resource) error {
|
func (rs *ResourceStore) create(r resource) error {
|
||||||
key, err := getKey(r)
|
key, err := getKey(r)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -90,7 +96,7 @@ func (rs *resourceStore) create(r resource) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
func (rs *resourceStore) read(rkey string, res resource) error {
|
func (rs *ResourceStore) read(rkey string, res resource) error {
|
||||||
typName, err := getTypeName(res)
|
typName, err := getTypeName(res)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -104,7 +110,7 @@ func (rs *resourceStore) read(rkey string, res resource) error {
|
|||||||
return rs.serializer.deserialize([]byte(resp.Node.Value), res)
|
return rs.serializer.deserialize([]byte(resp.Node.Value), res)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (rs *resourceStore) update(r resource) error {
|
func (rs *ResourceStore) update(r resource) error {
|
||||||
key, err := getKey(r)
|
key, err := getKey(r)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -120,13 +126,13 @@ func (rs *resourceStore) update(r resource) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
func (rs *resourceStore) delete(typename, rkey string) error {
|
func (rs *ResourceStore) delete(typename, rkey string) error {
|
||||||
key := typename + "/" + rkey
|
key := typename + "/" + rkey
|
||||||
_, err := rs.KeysAPI.Delete(context.Background(), key, nil) // ignore response
|
_, err := rs.KeysAPI.Delete(context.Background(), key, nil) // ignore response
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
func (rs *resourceStore) getAll(key string) ([]string, error) {
|
func (rs *ResourceStore) getAll(key string) ([]string, error) {
|
||||||
resp, err := rs.KeysAPI.Get(context.Background(), key, &client.GetOptions{Recursive: true})
|
resp, err := rs.KeysAPI.Get(context.Background(), key, &client.GetOptions{Recursive: true})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -139,10 +145,10 @@ func (rs *resourceStore) getAll(key string) ([]string, error) {
|
|||||||
return res, nil
|
return res, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (rs *resourceStore) writeFile(parentKey string, contents []byte) (string, string, error) {
|
func (rs *ResourceStore) writeFile(parentKey string, contents []byte) (string, string, error) {
|
||||||
uid := uuid.NewV4().String()
|
uid := uuid.NewV4().String()
|
||||||
|
|
||||||
err := rs.fileStore.write(uid, contents)
|
err := rs.FileStore.write(uid, contents)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", "", err
|
return "", "", err
|
||||||
}
|
}
|
||||||
@@ -150,14 +156,14 @@ func (rs *resourceStore) writeFile(parentKey string, contents []byte) (string, s
|
|||||||
parentKey = "file/" + parentKey
|
parentKey = "file/" + parentKey
|
||||||
resp, err := rs.KeysAPI.CreateInOrder(context.Background(), parentKey, uid, nil)
|
resp, err := rs.KeysAPI.CreateInOrder(context.Background(), parentKey, uid, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
_ = rs.fileStore.delete(uid)
|
_ = rs.FileStore.delete(uid)
|
||||||
return "", "", err
|
return "", "", err
|
||||||
}
|
}
|
||||||
|
|
||||||
return resp.Node.Key, uid, nil
|
return resp.Node.Key, uid, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (rs *resourceStore) readFile(key string, uid *string) ([]byte, error) {
|
func (rs *ResourceStore) readFile(key string, uid *string) ([]byte, error) {
|
||||||
key = "file/" + key
|
key = "file/" + key
|
||||||
resp, err := rs.KeysAPI.Get(context.Background(), key, &client.GetOptions{Sort: true})
|
resp, err := rs.KeysAPI.Get(context.Background(), key, &client.GetOptions{Sort: true})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -182,11 +188,11 @@ func (rs *resourceStore) readFile(key string, uid *string) ([]byte, error) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
contents, err := rs.fileStore.read(*uid)
|
contents, err := rs.FileStore.read(*uid)
|
||||||
return contents, err
|
return contents, err
|
||||||
}
|
}
|
||||||
|
|
||||||
func (rs *resourceStore) deleteFile(key string, uid string) error {
|
func (rs *ResourceStore) deleteFile(key string, uid string) error {
|
||||||
key = "file/" + key
|
key = "file/" + key
|
||||||
resp, err := rs.KeysAPI.Get(context.Background(), key, &client.GetOptions{Sort: true})
|
resp, err := rs.KeysAPI.Get(context.Background(), key, &client.GetOptions{Sort: true})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -204,7 +210,7 @@ func (rs *resourceStore) deleteFile(key string, uid string) error {
|
|||||||
return errors.New("won't delete unreferenced file")
|
return errors.New("won't delete unreferenced file")
|
||||||
}
|
}
|
||||||
|
|
||||||
err = rs.fileStore.delete(node.Value)
|
err = rs.FileStore.delete(node.Value)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -221,14 +227,14 @@ func (rs *resourceStore) deleteFile(key string, uid string) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (rs *resourceStore) deleteAllFiles(key string) error {
|
func (rs *ResourceStore) deleteAllFiles(key string) error {
|
||||||
key = "file/" + key
|
key = "file/" + key
|
||||||
resp, err := rs.KeysAPI.Get(context.Background(), key, &client.GetOptions{Sort: true})
|
resp, err := rs.KeysAPI.Get(context.Background(), key, &client.GetOptions{Sort: true})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
for _, u := range resp.Node.Nodes {
|
for _, u := range resp.Node.Nodes {
|
||||||
err = rs.fileStore.delete(u.Value)
|
err = rs.FileStore.delete(u.Value)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -47,19 +47,16 @@ func assert(b bool, msg string) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func getTestResourceStore() (*fileStore, client.KeysAPI, *resourceStore) {
|
func getTestResourceStore() (*FileStore, client.KeysAPI, *ResourceStore) {
|
||||||
// make a tmp dir
|
// make a tmp dir
|
||||||
dir, err := ioutil.TempDir("", "testFileStore")
|
dir, err := ioutil.TempDir("", "testFileStore")
|
||||||
panicIf(err)
|
panicIf(err)
|
||||||
fs := makeFileStore(dir)
|
fs := MakeFileStore(dir)
|
||||||
|
|
||||||
// assume etcd is running, connect to it
|
rs, err := MakeResourceStore(fs, []string{"http://localhost:2379"})
|
||||||
ks := getEtcdKeyAPI([]string{"http://localhost:2379"})
|
panicIf(err)
|
||||||
|
|
||||||
s := JsonSerializer{}
|
return fs, rs.KeysAPI, rs
|
||||||
rs := makeResourceStore(fs, ks, s)
|
|
||||||
|
|
||||||
return fs, ks, rs
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestResourceStore(t *testing.T) {
|
func TestResourceStore(t *testing.T) {
|
||||||
@@ -110,7 +107,7 @@ func TestResourceStore(t *testing.T) {
|
|||||||
assert(res[0] == tr, "value from retrieved list must equal updated value")
|
assert(res[0] == tr, "value from retrieved list must equal updated value")
|
||||||
|
|
||||||
// file tests
|
// file tests
|
||||||
fileKey := "resourceStoreTest"
|
fileKey := "ResourceStoreTest"
|
||||||
fileContents1 := []byte("hello")
|
fileContents1 := []byte("hello")
|
||||||
fileContents2 := []byte("world")
|
fileContents2 := []byte("world")
|
||||||
key, uid1, err := rs.writeFile(fileKey, fileContents1)
|
key, uid1, err := rs.writeFile(fileKey, fileContents1)
|
||||||
|
|||||||
@@ -27,3 +27,34 @@ func (e Error) Error() string {
|
|||||||
func MakeError(code int, msg string) Error {
|
func MakeError(code int, msg string) Error {
|
||||||
return Error{Code: errorCode(code), Message: msg}
|
return Error{Code: errorCode(code), Message: msg}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (err Error) HTTPStatus() int {
|
||||||
|
var code int
|
||||||
|
switch err.Code {
|
||||||
|
case ErrorNotFound:
|
||||||
|
code = 404
|
||||||
|
case ErrorInvalidArgument:
|
||||||
|
code = 400
|
||||||
|
case ErrorNoSpace:
|
||||||
|
code = 500
|
||||||
|
case ErrorNotAuthorized:
|
||||||
|
code = 403
|
||||||
|
default:
|
||||||
|
code = 500
|
||||||
|
}
|
||||||
|
return code
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetHTTPError(err error) (int, string) {
|
||||||
|
var msg string
|
||||||
|
var code int
|
||||||
|
fe, ok := err.(Error)
|
||||||
|
if ok {
|
||||||
|
msg = fe.Message
|
||||||
|
code = fe.HTTPStatus()
|
||||||
|
} else {
|
||||||
|
code = 500
|
||||||
|
msg = err.Error()
|
||||||
|
}
|
||||||
|
return code, msg
|
||||||
|
}
|
||||||
|
|||||||
@@ -10,8 +10,8 @@ const morgan = require('morgan');
|
|||||||
// Command line opts
|
// Command line opts
|
||||||
const argv = require('minimist')(process.argv.slice(1));
|
const argv = require('minimist')(process.argv.slice(1));
|
||||||
if (!argv.codepath) {
|
if (!argv.codepath) {
|
||||||
console.log("Codepath defaulting to /user.js");
|
argv.codepath = "/userfunc/user";
|
||||||
argv.codepath = "/user.js";
|
console.log("Codepath defaulting to ", argv.codepath);
|
||||||
}
|
}
|
||||||
if (!argv.port) {
|
if (!argv.port) {
|
||||||
console.log("Port defaulting to 8888");
|
console.log("Port defaulting to 8888");
|
||||||
@@ -79,7 +79,11 @@ app.all('/', function (req, res) {
|
|||||||
}
|
}
|
||||||
res.status(status).send(body);
|
res.status(status).send(body);
|
||||||
}
|
}
|
||||||
userFunction(context, callback);
|
try {
|
||||||
|
userFunction(context, callback);
|
||||||
|
} catch(e) {
|
||||||
|
callback(500, "Internal server error")
|
||||||
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
app.listen(argv.port);
|
app.listen(argv.port);
|
||||||
|
|||||||
@@ -42,10 +42,11 @@ func (fh *functionHandler) handler(responseWriter http.ResponseWriter, request *
|
|||||||
// Cache miss: request the Pool Manager to make a new service.
|
// Cache miss: request the Pool Manager to make a new service.
|
||||||
serviceUrl, poolErr := fh.getServiceForFunction()
|
serviceUrl, poolErr := fh.getServiceForFunction()
|
||||||
if poolErr != nil {
|
if poolErr != nil {
|
||||||
// now we're really screwed
|
|
||||||
log.Printf("Failed to get service for function (%v,%v): %v",
|
log.Printf("Failed to get service for function (%v,%v): %v",
|
||||||
fh.Function.Name, fh.Function.Uid, poolErr)
|
fh.Function.Name, fh.Function.Uid, poolErr)
|
||||||
responseWriter.WriteHeader(500) // TODO: make this smarter based on the actual error
|
// We might want a specific error code or header for fission
|
||||||
|
// failures as opposed to user function bugs.
|
||||||
|
http.Error(responseWriter, poolErr.Error(), 500)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+4
-27
@@ -42,19 +42,9 @@ package router
|
|||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"github.com/gorilla/mux"
|
"github.com/gorilla/mux"
|
||||||
flag "github.com/ogier/pflag"
|
|
||||||
"net/http"
|
"net/http"
|
||||||
)
|
)
|
||||||
|
|
||||||
type (
|
|
||||||
options struct {
|
|
||||||
port int
|
|
||||||
poolManagerUrl string
|
|
||||||
controllerUrl string
|
|
||||||
//...
|
|
||||||
}
|
|
||||||
)
|
|
||||||
|
|
||||||
// request url ---[mux]---> Function(name,uid) ----[fmap]----> k8s service url
|
// request url ---[mux]---> Function(name,uid) ----[fmap]----> k8s service url
|
||||||
|
|
||||||
// request url ---[trigger]---> Function(name, deployment) ----[deployment]----> Function(name, uid) ----[pool mgr]---> k8s service url
|
// request url ---[trigger]---> Function(name, deployment) ----[deployment]----> Function(name, uid) ----[pool mgr]---> k8s service url
|
||||||
@@ -66,27 +56,14 @@ func router(httpTriggerSet *HTTPTriggerSet) *mutableRouter {
|
|||||||
return mr
|
return mr
|
||||||
}
|
}
|
||||||
|
|
||||||
func server(port int, httpTriggerSet *HTTPTriggerSet) {
|
func serve(port int, httpTriggerSet *HTTPTriggerSet) {
|
||||||
mr := router(httpTriggerSet)
|
mr := router(httpTriggerSet)
|
||||||
url := fmt.Sprintf(":%v", port)
|
url := fmt.Sprintf(":%v", port)
|
||||||
http.ListenAndServe(url, mr)
|
http.ListenAndServe(url, mr)
|
||||||
}
|
}
|
||||||
|
|
||||||
func getOptions() *options {
|
func Start(port int, controllerUrl string, poolmgrUrl string) {
|
||||||
options := &options{}
|
|
||||||
|
|
||||||
flag.IntVar(&options.port, "port", 80, "Port to listen on")
|
|
||||||
|
|
||||||
// default to using dns service discovery
|
|
||||||
flag.StringVar(&options.poolManagerUrl, "poolmanager_url", "http://poolmanager/", "URL for the PoolManager service")
|
|
||||||
flag.StringVar(&options.controllerUrl, "controller_url", "http://controller/", "URL for the controller service")
|
|
||||||
|
|
||||||
return options
|
|
||||||
}
|
|
||||||
|
|
||||||
func main() {
|
|
||||||
options := getOptions()
|
|
||||||
fmap := makeFunctionServiceMap()
|
fmap := makeFunctionServiceMap()
|
||||||
triggers := makeHTTPTriggerSet(fmap, options.controllerUrl, options.poolManagerUrl)
|
triggers := makeHTTPTriggerSet(fmap, controllerUrl, poolmgrUrl)
|
||||||
server(options.port, triggers)
|
serve(port, triggers)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -38,7 +38,7 @@ func TestRouter(t *testing.T) {
|
|||||||
triggers.triggers = append(triggers.triggers, fission.HTTPTrigger{UrlPattern: triggerUrl, Function: *fn})
|
triggers.triggers = append(triggers.triggers, fission.HTTPTrigger{UrlPattern: triggerUrl, Function: *fn})
|
||||||
|
|
||||||
port := 4242
|
port := 4242
|
||||||
go server(port, triggers)
|
go serve(port, triggers)
|
||||||
time.Sleep(100 * time.Millisecond)
|
time.Sleep(100 * time.Millisecond)
|
||||||
|
|
||||||
testUrl := fmt.Sprintf("http://localhost:%v%v", port, triggerUrl)
|
testUrl := fmt.Sprintf("http://localhost:%v%v", port, triggerUrl)
|
||||||
|
|||||||
Reference in New Issue
Block a user