diff --git a/Documentation/Architecture.md b/Documentation/Architecture.md index 8b65469b..aec937b6 100644 --- a/Documentation/Architecture.md +++ b/Documentation/Architecture.md @@ -95,6 +95,19 @@ While a few simple retries are done, there isn't yet a reliable message bus between Kubewatcher and the function. Work for this is tracked in issue #64. +Message Queue Trigger +--------------------- + +A message queue trigger binds a message queue topic to a function: +events from that topic cause the function to be invoked with the +message as the body of the request. The trigger may also contain a +response topic: if specified, the function's output is sent to this +response. + +Here's a diagram of the components: + +![Message queue trigger Diagram](https://user-images.githubusercontent.com/202578/27012344-9457cb24-4f00-11e7-8d6b-926ff01637b3.jpg) + Environment Container --------------------- @@ -113,7 +126,7 @@ volume shared between fetcher and this environment container. Poolmgr then requests the container to load the function. Logger ------------ +------ Logger helps to forward function logs to centralized db service for log persistence. Currently only influxdb is supported to store logs. diff --git a/INSTALL.md b/INSTALL.md index f73dbce9..eb3b1ae6 100644 --- a/INSTALL.md +++ b/INSTALL.md @@ -200,3 +200,22 @@ To setup Fission-ui with fission in k8s is simple: Then open `http://node-ip:31319` to use Fission-ui. For more infomation, please check out [Fission-ui Readme](https://github.com/fission/fission-ui/blob/master/README.md). + +### Install NATS for message-queue based triggers (Optional) + +Fission supports message queue triggers that allow you to invoke +functions based on events in a queue. For now, NATS-Streaming is the +only supported message queue. + +You can install NATS Streaming on your Kubernetes cluster with: + +``` + $ kubectl create -f fission-nats.yaml +``` + +You can subscribe to a NATS Streaming queue with a command like this: +(See `fission mqtrigger --help` for details) + +``` + $ fission mqtrigger create --name myQueueTrigger --function processEvent --topic "myQueue.request" +``` diff --git a/controller/api.go b/controller/api.go index 929d401c..3504ba72 100644 --- a/controller/api.go +++ b/controller/api.go @@ -36,6 +36,7 @@ type ( FunctionStore HTTPTriggerStore TimeTriggerStore + MessageQueueTriggerStore EnvironmentStore WatchStore } @@ -49,11 +50,12 @@ type ( func MakeAPI(rs *ResourceStore) *API { api := &API{ - FunctionStore: FunctionStore{ResourceStore: *rs}, - HTTPTriggerStore: HTTPTriggerStore{ResourceStore: *rs}, - TimeTriggerStore: TimeTriggerStore{ResourceStore: *rs}, - EnvironmentStore: EnvironmentStore{ResourceStore: *rs}, - WatchStore: WatchStore{ResourceStore: *rs}, + FunctionStore: FunctionStore{ResourceStore: *rs}, + HTTPTriggerStore: HTTPTriggerStore{ResourceStore: *rs}, + TimeTriggerStore: TimeTriggerStore{ResourceStore: *rs}, + MessageQueueTriggerStore: MessageQueueTriggerStore{ResourceStore: *rs}, + EnvironmentStore: EnvironmentStore{ResourceStore: *rs}, + WatchStore: WatchStore{ResourceStore: *rs}, } return api } @@ -130,6 +132,12 @@ func (api *API) Serve(port int) { r.HandleFunc("/v1/triggers/time/{timeTrigger}", api.TimeTriggerApiUpdate).Methods("PUT") r.HandleFunc("/v1/triggers/time/{timeTrigger}", api.TimeTriggerApiDelete).Methods("DELETE") + r.HandleFunc("/v1/triggers/messagequeue", api.MessageQueueTriggerApiList).Methods("GET") + r.HandleFunc("/v1/triggers/messagequeue", api.MessageQueueApiCreate).Methods("POST") + r.HandleFunc("/v1/triggers/messagequeue/{mqTrigger}", api.MessageQueueApiGet).Methods("GET") + r.HandleFunc("/v1/triggers/messagequeue/{mqTrigger}", api.MessageQueueApiUpdate).Methods("PUT") + r.HandleFunc("/v1/triggers/messagequeue/{mqTrigger}", api.MessageQueueApiDelete).Methods("DELETE") + r.HandleFunc("/proxy/{dbType}", api.FunctionLogsApiPost).Methods("POST") address := fmt.Sprintf(":%v", port) diff --git a/controller/client/client.go b/controller/client/client.go index 3fa15a3d..571ee086 100644 --- a/controller/client/client.go +++ b/controller/client/client.go @@ -640,3 +640,116 @@ func (c *Client) TimeTriggerList() ([]fission.TimeTrigger, error) { return triggers, nil } + +func (c *Client) MessageQueueTriggerCreate(t *fission.MessageQueueTrigger) (*fission.Metadata, error) { + reqbody, err := json.Marshal(t) + if err != nil { + return nil, err + } + + resp, err := http.Post(c.url("triggers/messagequeue"), "application/json", bytes.NewReader(reqbody)) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + body, err := c.handleCreateResponse(resp) + if err != nil { + return nil, err + } + + var m fission.Metadata + err = json.Unmarshal(body, &m) + if err != nil { + return nil, err + } + + return &m, nil +} + +func (c *Client) MessageQueueTriggerGet(m *fission.Metadata) (*fission.MessageQueueTrigger, error) { + relativeUrl := fmt.Sprintf("triggers/messagequeue/%v", m.Name) + if len(m.Uid) > 0 { + relativeUrl += fmt.Sprintf("?uid=%v", m.Uid) + } + + resp, err := http.Get(c.url(relativeUrl)) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + body, err := c.handleResponse(resp) + if err != nil { + return nil, err + } + + var t fission.MessageQueueTrigger + err = json.Unmarshal(body, &t) + if err != nil { + return nil, err + } + + return &t, nil +} + +func (c *Client) MessageQueueTriggerUpdate(mqTrigger *fission.MessageQueueTrigger) (*fission.Metadata, error) { + reqbody, err := json.Marshal(mqTrigger) + if err != nil { + return nil, err + } + relativeUrl := fmt.Sprintf("triggers/messagequeue/%v", mqTrigger.Metadata.Name) + + resp, err := c.put(relativeUrl, "application/json", reqbody) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + body, err := c.handleResponse(resp) + if err != nil { + return nil, err + } + + var m fission.Metadata + err = json.Unmarshal(body, &m) + if err != nil { + return nil, err + } + return &m, nil +} + +func (c *Client) MessageQueueTriggerDelete(m *fission.Metadata) error { + relativeUrl := fmt.Sprintf("triggers/messagequeue/%v", m.Name) + if len(m.Uid) > 0 { + relativeUrl += fmt.Sprintf("?uid=%v", m.Uid) + } + err := c.delete(relativeUrl) + return err +} + +func (c *Client) MessageQueueTriggerList(mqType string) ([]fission.MessageQueueTrigger, error) { + relativeUrl := "triggers/messagequeue" + if len(mqType) > 0 { + relativeUrl += fmt.Sprintf("?mqtype=%v", mqType) + } + + resp, err := http.Get(c.url(relativeUrl)) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + body, err := c.handleResponse(resp) + if err != nil { + return nil, err + } + + triggers := make([]fission.MessageQueueTrigger, 0) + err = json.Unmarshal(body, &triggers) + if err != nil { + return nil, err + } + + return triggers, nil +} diff --git a/controller/mqTriggerApi.go b/controller/mqTriggerApi.go new file mode 100644 index 00000000..91fd9ea2 --- /dev/null +++ b/controller/mqTriggerApi.go @@ -0,0 +1,166 @@ +/* +Copyright 2017 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package controller + +import ( + "encoding/json" + "io/ioutil" + "net/http" + + "github.com/gorilla/mux" + log "github.com/sirupsen/logrus" + + "github.com/fission/fission" +) + +func (api *API) MessageQueueTriggerApiList(w http.ResponseWriter, r *http.Request) { + mqType := r.FormValue("mqtype") + triggers, err := api.MessageQueueTriggerStore.List(mqType) + if err != nil { + api.respondWithError(w, err) + return + } + resp, err := json.Marshal(triggers) + if err != nil { + api.respondWithError(w, err) + return + } + api.respondWithSuccess(w, resp) +} + +func (api *API) MessageQueueApiCreate(w http.ResponseWriter, r *http.Request) { + body, err := ioutil.ReadAll(r.Body) + if err != nil { + api.respondWithError(w, err) + return + } + + var mqTrigger fission.MessageQueueTrigger + err = json.Unmarshal(body, &mqTrigger) + if err != nil { + api.respondWithError(w, err) + return + } + + // trigger name must not conflict with any other trigger + // even they are different message queue type + triggers, err := api.MessageQueueTriggerStore.List("") + if err != nil { + api.respondWithError(w, err) + return + } + for _, trigger := range triggers { + if trigger.Name == mqTrigger.Name { + err = fission.MakeError(fission.ErrorNameExists, + "Message queue trigger with same name already exists") + api.respondWithError(w, err) + return + } + } + + // save trigger info + uid, err := api.MessageQueueTriggerStore.Create(&mqTrigger) + if err != nil { + api.respondWithError(w, err) + return + } + + mqTriggerMeta := fission.Metadata{Name: mqTrigger.Metadata.Name, Uid: uid} + resp, err := json.Marshal(mqTriggerMeta) + if err != nil { + api.respondWithError(w, err) + return + } + w.WriteHeader(http.StatusCreated) + api.respondWithSuccess(w, resp) +} + +func (api *API) MessageQueueApiGet(w http.ResponseWriter, r *http.Request) { + vars := mux.Vars(r) + mqTriggerMeta := fission.Metadata{ + Name: vars["mqTrigger"], + Uid: r.FormValue("uid"), // empty if uid is absent + } + mqTrigger, err := api.MessageQueueTriggerStore.Get(&mqTriggerMeta) + if err != nil { + api.respondWithError(w, err) + return + } + resp, err := json.Marshal(mqTrigger) + if err != nil { + api.respondWithError(w, err) + return + } + api.respondWithSuccess(w, resp) +} + +func (api *API) MessageQueueApiUpdate(w http.ResponseWriter, r *http.Request) { + vars := mux.Vars(r) + mqtName := vars["mqTrigger"] + + body, err := ioutil.ReadAll(r.Body) + if err != nil { + api.respondWithError(w, err) + return + } + + var mqTrigger fission.MessageQueueTrigger + err = json.Unmarshal(body, &mqTrigger) + if err != nil { + api.respondWithError(w, err) + return + } + + if mqtName != mqTrigger.Metadata.Name { + err = fission.MakeError(fission.ErrorInvalidArgument, "Message queue trigger name doesn't match URL") + api.respondWithError(w, err) + return + } + + uid, err := api.MessageQueueTriggerStore.Update(&mqTrigger) + if err != nil { + api.respondWithError(w, err) + return + } + + mqTriggerMeta := fission.Metadata{Name: mqTrigger.Metadata.Name, Uid: uid} + resp, err := json.Marshal(mqTriggerMeta) + if err != nil { + api.respondWithError(w, err) + return + } + api.respondWithSuccess(w, resp) +} + +func (api *API) MessageQueueApiDelete(w http.ResponseWriter, r *http.Request) { + vars := mux.Vars(r) + mqTriggerMeta := fission.Metadata{ + Name: vars["mqTrigger"], + Uid: r.FormValue("uid"), // empty if uid is absent + } + + if len(mqTriggerMeta.Uid) == 0 { + log.WithFields(log.Fields{"mqTrigger": mqTriggerMeta.Name}).Info("Deleting all versions") + } + + err := api.MessageQueueTriggerStore.Delete(mqTriggerMeta) + if err != nil { + api.respondWithError(w, err) + return + } + api.respondWithSuccess(w, []byte("")) +} diff --git a/controller/mqTriggerStore.go b/controller/mqTriggerStore.go new file mode 100644 index 00000000..80fb6bf2 --- /dev/null +++ b/controller/mqTriggerStore.go @@ -0,0 +1,80 @@ +/* +CopyrigmqTrigger 2017 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + mqTriggertp://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package controller + +import ( + "github.com/fission/fission" + uuid "github.com/satori/go.uuid" +) + +type MessageQueueTriggerStore struct { + ResourceStore +} + +func (mqs *MessageQueueTriggerStore) Create(mqTrigger *fission.MessageQueueTrigger) (string, error) { + mqTrigger.Metadata.Uid = uuid.NewV4().String() + return mqTrigger.Metadata.Uid, mqs.ResourceStore.create(mqTrigger) +} + +func (mqs *MessageQueueTriggerStore) Get(m *fission.Metadata) (*fission.MessageQueueTrigger, error) { + var mqTrigger fission.MessageQueueTrigger + err := mqs.ResourceStore.read(m.Name, &mqTrigger) + if err != nil { + return nil, err + } + return &mqTrigger, nil +} + +func (mqs *MessageQueueTriggerStore) Update(mqTrigger *fission.MessageQueueTrigger) (string, error) { + mqTrigger.Metadata.Uid = uuid.NewV4().String() + return mqTrigger.Metadata.Uid, mqs.ResourceStore.update(mqTrigger) +} + +func (mqs *MessageQueueTriggerStore) Delete(m fission.Metadata) error { + typeName, err := getTypeName(fission.MessageQueueTrigger{}) + if err != nil { + return err + } + return mqs.ResourceStore.delete(typeName, m.Name) +} + +func (mqs *MessageQueueTriggerStore) List(mqType string) ([]fission.MessageQueueTrigger, error) { + typeName, err := getTypeName(fission.MessageQueueTrigger{}) + if err != nil { + return nil, err + } + bufs, err := mqs.ResourceStore.getAll(typeName) + if err != nil { + return nil, err + } + + triggers := make([]fission.MessageQueueTrigger, 0, len(bufs)) + js := JsonSerializer{} + for _, buf := range bufs { + var mqTrigger fission.MessageQueueTrigger + err = js.deserialize([]byte(buf), &mqTrigger) + if err != nil { + return nil, err + } + if len(mqType) > 0 && mqType != mqTrigger.MessageQueueType { + continue + } + triggers = append(triggers, mqTrigger) + } + + return triggers, nil +} diff --git a/fission-bundle/main.go b/fission-bundle/main.go index 56282cde..30e295c4 100644 --- a/fission-bundle/main.go +++ b/fission-bundle/main.go @@ -8,6 +8,7 @@ import ( "github.com/fission/fission/controller" "github.com/fission/fission/kubewatcher" "github.com/fission/fission/logger" + "github.com/fission/fission/mqtrigger" "github.com/fission/fission/poolmgr" "github.com/fission/fission/router" "github.com/fission/fission/timer" @@ -61,6 +62,13 @@ func runTimer(controllerUrl, routerUrl string) { } } +func runMessageQueueMgr(controllerUrl, routerUrl string) { + err := messagequeue.Start(controllerUrl, routerUrl) + if err != nil { + log.Fatalf("Error starting timer: %v", err) + } +} + func getPort(portArg interface{}) int { portArgStr := portArg.(string) port, err := strconv.Atoi(portArgStr) @@ -96,6 +104,7 @@ Usage: fission-bundle --kubewatcher [--controllerUrl= --routerUrl=] fission-bundle --logger fission-bundle --timer [--controllerUrl= --routerUrl=] + fission-bundle --mqt [--controllerUrl= --routerUrl=] Options: --controllerPort= Port that the controller should listen on. --routerPort= Port that the router should listen on. @@ -109,6 +118,7 @@ Options: --kubewatcher Start Kubernetes events watcher. --logger Start logger. --timer Start Timer. + --mqt Start message queue trigger. ` arguments, err := docopt.Parse(usage, nil, true, "fission-bundle", false) if err != nil { @@ -149,5 +159,9 @@ Options: runTimer(controllerUrl, routerUrl) } + if arguments["--mqt"] == true { + runMessageQueueMgr(controllerUrl, routerUrl) + } + select {} } diff --git a/fission-cloud.yaml b/fission-cloud.yaml index 0b58433c..cf34d847 100644 --- a/fission-cloud.yaml +++ b/fission-cloud.yaml @@ -48,4 +48,22 @@ spec: targetPort: 8086 nodePort: 31315 selector: - svc: influxdb \ No newline at end of file + svc: influxdb + +--- + +apiVersion: v1 +kind: Service +metadata: + name: nats-streaming + namespace: fission + labels: + svc: nats-streaming +spec: + type: LoadBalancer + ports: + - port: 4222 + targetPort: 4222 + nodePort: 31316 + selector: + svc: nats-streaming \ No newline at end of file diff --git a/fission-nats.yaml b/fission-nats.yaml new file mode 100644 index 00000000..9e0f5eeb --- /dev/null +++ b/fission-nats.yaml @@ -0,0 +1,60 @@ +# To enable nats authentication, please follow the instruction described in +# http://nats.io/documentation/server/gnatsd-authentication/. +# And dont forget to change the MESSAGE_QUEUE_URL in mqtrigger deployment. + +apiVersion: extensions/v1beta1 +kind: Deployment +metadata: + name: mqtrigger + namespace: fission +spec: + replicas: 1 + template: + metadata: + labels: + svc: mqtrigger + spec: + containers: + - name: mqtrigger + image: fission/fission-bundle + command: ["/fission-bundle"] + args: ["--mqt"] + env: + - name: MESSAGE_QUEUE_TYPE + value: nats-streaming + - name: MESSAGE_QUEUE_URL + value: nats://nats-streaming:4222 + +--- +apiVersion: extensions/v1beta1 +kind: Deployment +metadata: + labels: + svc: nats-streaming + name: nats-streaming + namespace: fission +spec: + replicas: 1 + selector: + matchLabels: + svc: nats-streaming + strategy: + rollingUpdate: + maxSurge: 1 + maxUnavailable: 1 + type: RollingUpdate + template: + metadata: + creationTimestamp: null + labels: + svc: nats-streaming + spec: + containers: + - name: nats-streaming + image: nats-streaming + args: ["--cluster_id", "fissionMQTrigger"] + ports: + - containerPort: 4222 + hostPort: 4222 + protocol: TCP + diff --git a/fission-nodeport.yaml b/fission-nodeport.yaml index ea9fcbc9..62730fc9 100644 --- a/fission-nodeport.yaml +++ b/fission-nodeport.yaml @@ -48,4 +48,22 @@ spec: targetPort: 8086 nodePort: 31315 selector: - svc: influxdb \ No newline at end of file + svc: influxdb + +--- + +apiVersion: v1 +kind: Service +metadata: + name: nats-streaming + namespace: fission + labels: + svc: nats-streaming +spec: + type: NodePort + ports: + - port: 4222 + targetPort: 4222 + nodePort: 31316 + selector: + svc: nats-streaming \ No newline at end of file diff --git a/fission/main.go b/fission/main.go index 301b827d..732e96a7 100644 --- a/fission/main.go +++ b/fission/main.go @@ -82,6 +82,21 @@ func main() { {Name: "list", Usage: "List Time triggers", Flags: []cli.Flag{}, Action: ttList}, } + // Message queue trigger + mqtNameFlag := cli.StringFlag{Name: "name", Usage: "Message queue Trigger name"} + mqtFnNameFlag := cli.StringFlag{Name: "function", Usage: "Function name"} + mqtFnUidFlag := cli.StringFlag{Name: "uid", Usage: "Function UID (optional; uses latest if unspecified)"} + mqtMQTypeFlag := cli.StringFlag{Name: "mqtype", Usage: "Message queue type, e.g. nats-streaming (optional; uses \"nats-streaming\" if unspecified)"} + mqtTopicFlag := cli.StringFlag{Name: "topic", Usage: "Message queue Topic the trigger listens on"} + mqtRespTopicFlag := cli.StringFlag{Name: "resptopic", Usage: "Topic that the function response is sent on (optional; response discarded if unspecified)"} + mqtSubcommands := []cli.Command{ + {Name: "create", Aliases: []string{"add"}, Usage: "Create Message queue trigger", Flags: []cli.Flag{mqtNameFlag, mqtFnNameFlag, mqtFnUidFlag, mqtMQTypeFlag, mqtTopicFlag, mqtRespTopicFlag}, Action: mqtCreate}, + {Name: "get", Usage: "Get message queue trigger", Flags: []cli.Flag{}, Action: mqtGet}, + {Name: "update", Usage: "Update message queue trigger", Flags: []cli.Flag{mqtNameFlag, mqtTopicFlag, mqtRespTopicFlag}, Action: mqtUpdate}, + {Name: "delete", Usage: "Delete message queue trigger", Flags: []cli.Flag{mqtNameFlag}, Action: mqtDelete}, + {Name: "list", Usage: "List message queue triggers", Flags: []cli.Flag{mqtMQTypeFlag}, Action: mqtList}, + } + // environments envNameFlag := cli.StringFlag{Name: "name", Usage: "Environment name"} envImageFlag := cli.StringFlag{Name: "image", Usage: "Environment image URL"} @@ -112,6 +127,7 @@ func main() { {Name: "function", Aliases: []string{"fn"}, Usage: "Create, update and manage functions", Subcommands: fnSubcommands}, {Name: "httptrigger", Aliases: []string{"ht", "route"}, Usage: "Manage HTTP triggers (routes) for functions", Subcommands: htSubcommands}, {Name: "timetrigger", Aliases: []string{"tt", "timer"}, Usage: "Manage Time triggers (timers) for functions", Subcommands: ttSubcommands}, + {Name: "mqtrigger", Aliases: []string{"mqt", "messagequeue"}, Usage: "Manage message queue triggers for functions", Subcommands: mqtSubcommands}, {Name: "environment", Aliases: []string{"env"}, Usage: "Manage environments", Subcommands: envSubcommands}, {Name: "watch", Aliases: []string{"w"}, Usage: "Manage watches", Subcommands: wSubCommands}, diff --git a/fission/mqtrigger.go b/fission/mqtrigger.go new file mode 100644 index 00000000..d0e7b4d8 --- /dev/null +++ b/fission/mqtrigger.go @@ -0,0 +1,155 @@ +/* +Copyrigtt 2017 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + tttp://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package main + +import ( + "fmt" + "os" + "text/tabwriter" + + "github.com/satori/go.uuid" + "github.com/urfave/cli" + + "github.com/fission/fission" + "github.com/fission/fission/mqtrigger/messageQueue" +) + +func mqtCreate(c *cli.Context) error { + client := getClient(c.GlobalString("server")) + + mqtName := c.String("name") + if len(mqtName) == 0 { + mqtName = uuid.NewV4().String() + } + fnName := c.String("function") + if len(fnName) == 0 { + fatal("Need a function name to create a trigger, use --function") + } + fnUid := c.String("uid") + mqType := c.String("mqtype") + switch mqType { + case "": + mqType = messageQueue.NATS + case messageQueue.NATS: + mqType = messageQueue.NATS + default: + fatal("Unknown message queue type, currently only \"nats-streaming\" is supported") + } + + // TODO: check topic availability + topic := c.String("topic") + if len(topic) == 0 { + fatal("Listen topic cannot be empty") + } + respTopic := c.String("resptopic") + + if topic == respTopic { + fatal("Listen topic should not equal to response topic") + } + + checkMQTopicAvailability(mqType, topic, respTopic) + + fnMeta := fission.Metadata{ + Name: fnName, + Uid: fnUid, + } + + mqt := fission.MessageQueueTrigger{ + Metadata: fission.Metadata{ + Name: mqtName, + }, + Function: fnMeta, + MessageQueueType: mqType, + Topic: topic, + ResponseTopic: respTopic, + } + + _, err := client.MessageQueueTriggerCreate(&mqt) + checkErr(err, "create message queue trigger") + + fmt.Printf("trigger '%s' created\n", mqtName) + return err +} + +func mqtGet(c *cli.Context) error { + return nil +} + +func mqtUpdate(c *cli.Context) error { + client := getClient(c.GlobalString("server")) + mqtName := c.String("name") + if len(mqtName) == 0 { + fatal("Need name of trigger, use --name") + } + topic := c.String("topic") + respTopic := c.String("resptopic") + + mqt, err := client.MessageQueueTriggerGet(&fission.Metadata{Name: mqtName}) + checkErr(err, "get Time trigger") + + checkMQTopicAvailability(mqt.MessageQueueType, topic, respTopic) + + mqt.Topic = topic + mqt.ResponseTopic = respTopic + + _, err = client.MessageQueueTriggerUpdate(mqt) + checkErr(err, "update Time trigger") + + fmt.Printf("trigger '%v' updated\n", mqtName) + return nil +} + +func mqtDelete(c *cli.Context) error { + client := getClient(c.GlobalString("server")) + mqtName := c.String("name") + if len(mqtName) == 0 { + fatal("Need name of trigger to delete, use --name") + } + + err := client.MessageQueueTriggerDelete(&fission.Metadata{Name: mqtName}) + checkErr(err, "delete trigger") + + fmt.Printf("trigger '%v' deleted\n", mqtName) + return nil +} + +func mqtList(c *cli.Context) error { + client := getClient(c.GlobalString("server")) + + mqts, err := client.MessageQueueTriggerList(c.String("mqtype")) + checkErr(err, "list message queue triggers") + + w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) + + fmt.Fprintf(w, "%v\t%v\t%v\t%v\t%v\t%v\n", + "NAME", "FUNCTION_NAME", "FUNCTION_UID", "MESSAGE_QUEUE_TYPE", "TOPIC", "RESPONSE_TOPIC") + for _, mqt := range mqts { + fmt.Fprintf(w, "%v\t%v\t%v\t%v\t%v\t%v\n", + mqt.Metadata.Name, mqt.Function.Name, mqt.Function.Uid, mqt.MessageQueueType, mqt.Topic, mqt.ResponseTopic) + } + w.Flush() + + return nil +} + +func checkMQTopicAvailability(mqType string, topics ...string) { + for _, t := range topics { + if len(t) > 0 && !messageQueue.IsTopicValid(mqType, t) { + fatal(fmt.Sprintf("Invalid topic for %s: %s", mqType, t)) + } + } +} diff --git a/glide.lock b/glide.lock index 4d9e37b3..347e05dc 100644 --- a/glide.lock +++ b/glide.lock @@ -88,6 +88,21 @@ imports: - buffer - jlexer - jwriter +- name: github.com/nats-io/go-nats + version: 6949c8e06a246e4177961aab22940b5c411e48f0 + subpackages: + - encoders/builtin + - util +- name: github.com/nats-io/go-nats-streaming + version: 6e620057a207bd61e992c1c5b6a2de7b6a4cb010 + subpackages: + - pb +- name: github.com/nats-io/nats-streaming-server + version: 7a922646104d066c98527959455993231fadd168 + subpackages: + - util +- name: github.com/nats-io/nuid + version: 289cccf02c178dc782430d534e3c1f5b72af807f - name: github.com/pborman/uuid version: 3d4f2ba23642d3cfd06bd4b54cf03d99d95c0f1b - name: github.com/PuerkitoBio/purell diff --git a/glide.yaml b/glide.yaml index 64c5bb82..9c324a89 100644 --- a/glide.yaml +++ b/glide.yaml @@ -34,3 +34,7 @@ import: subpackages: - client/v2 - package: github.com/robfig/cron +- package: github.com/nats-io/go-nats-streaming + version: ^v0.3.4 +- package: github.com/nats-io/nats-streaming-server + version: ^v0.4.0 diff --git a/hack/release.sh b/hack/release.sh index c70c2a18..81366f5a 100755 --- a/hack/release.sh +++ b/hack/release.sh @@ -93,6 +93,7 @@ build_yaml() { cat fission-logger.yaml | sed "s#fission/fission-bundle#$tag#g" > $outdir/fission-logger.yaml cat fission-openshift.yaml | sed "s#fission/fission-bundle#$tag#g" > $outdir/fission-openshift.yaml cat fission-rbac.yaml | sed "s#fission/fission-bundle#$tag#g" > $outdir/fission-rbac.yaml + cat fission-nats.yaml | sed "s#fission/fission-bundle#$tag#g" > $outdir/fission-nats.yaml cp fission-nodeport.yaml $outdir cp fission-cloud.yaml $outdir @@ -138,7 +139,6 @@ make_github_release() { --tag $gittag \ --name "Nightly release for $(date +%Y-%b-%d)" \ --description "Nightly release for $(date +%Y-%b-%d)" \ - --pre-release # attach files @@ -165,7 +165,7 @@ make_github_release() { --file $BUILDDIR/cli/windows/fission.exe # yamls - yaml_files="fission.yaml fission-logger.yaml fission-rbac.yaml fission-openshift.yaml fission-nodeport.yaml fission-cloud.yaml" + yaml_files="fission.yaml fission-logger.yaml fission-rbac.yaml fission-openshift.yaml fission-nodeport.yaml fission-cloud.yaml fission-nats.yaml" for f in $yaml_files do gothub upload \ diff --git a/mqtrigger/main.go b/mqtrigger/main.go new file mode 100644 index 00000000..31c0a9f2 --- /dev/null +++ b/mqtrigger/main.go @@ -0,0 +1,38 @@ +/* +Copyright 2016 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package messagequeue + +import ( + "os" + + controllerClient "github.com/fission/fission/controller/client" + "github.com/fission/fission/mqtrigger/messageQueue" +) + +func Start(controllerUrl string, routerUrl string) error { + controller := controllerClient.MakeClient(controllerUrl) + + // Message queue type: nats is the only supported one for now + mqType := os.Getenv("MESSAGE_QUEUE_TYPE") + mqUrl := os.Getenv("MESSAGE_QUEUE_URL") + mqCfg := messageQueue.MessageQueueConfig{ + MQType: mqType, + Url: mqUrl, + } + messageQueue.MakeMessageQueueTriggerManager(controller, routerUrl, mqCfg) + return nil +} diff --git a/mqtrigger/messageQueue/messageQueue.go b/mqtrigger/messageQueue/messageQueue.go new file mode 100644 index 00000000..14356b9c --- /dev/null +++ b/mqtrigger/messageQueue/messageQueue.go @@ -0,0 +1,226 @@ +/* +Copyright 2016 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package messageQueue + +import ( + "errors" + "time" + + log "github.com/sirupsen/logrus" + + "github.com/fission/fission" + controllerClient "github.com/fission/fission/controller/client" +) + +const ( + NATS string = "nats-streaming" +) + +const ( + ADD_TRIGGER requestType = iota + DELETE_TRIGGER + GET_ALL_TRIGGERS +) + +type ( + messageQueueSubscription interface{} + + requestType int + + MessageQueueConfig struct { + MQType string + Url string + } + + MessageQueue interface { + subscribe(trigger fission.MessageQueueTrigger) (messageQueueSubscription, error) + unsubscribe(triggerSub messageQueueSubscription) error + } + + MessageQueueTriggerManager struct { + reqChan chan request + mqCfg MessageQueueConfig + triggers map[string]*triggerSubscription + controller *controllerClient.Client + messageQueue MessageQueue + } + + triggerSubscription struct { + fission.Metadata + funcMeta fission.Metadata + subscription messageQueueSubscription + } + + request struct { + requestType + triggerSub *triggerSubscription + respChan chan response + } + response struct { + err error + triggers *map[string]messageQueueSubscription + } +) + +func MakeMessageQueueTriggerManager(ctrlClient *controllerClient.Client, + routerUrl string, mqConfig MessageQueueConfig) *MessageQueueTriggerManager { + + var messageQueue MessageQueue + var err error + + mqTriggerMgr := MessageQueueTriggerManager{ + reqChan: make(chan request), + triggers: make(map[string]*triggerSubscription), + controller: ctrlClient, + } + switch mqConfig.MQType { + case NATS: + messageQueue, err = makeNatsMessageQueue(routerUrl, mqConfig) + default: + err = errors.New("No matched message queue type found") + } + if err != nil { + log.Fatalf("Failed to connect to remote message queue server: %v", err) + } + mqTriggerMgr.messageQueue = messageQueue + go mqTriggerMgr.service() + go mqTriggerMgr.syncTriggers() + return &mqTriggerMgr +} + +func (mqt *MessageQueueTriggerManager) service() { + for { + req := <-mqt.reqChan + switch req.requestType { + case ADD_TRIGGER: + var err error + triggerUid := req.triggerSub.Uid + if _, ok := mqt.triggers[triggerUid]; ok { + err = errors.New("Trigger already exists") + } else { + mqt.triggers[triggerUid] = req.triggerSub + } + req.respChan <- response{err: err} + case GET_ALL_TRIGGERS: + copyTriggers := make(map[string]messageQueueSubscription) + for key, val := range mqt.triggers { + copyTriggers[key] = val + } + req.respChan <- response{triggers: ©Triggers} + case DELETE_TRIGGER: + triggerUid := req.triggerSub.Uid + delete(mqt.triggers, triggerUid) + } + } +} + +func (mqt *MessageQueueTriggerManager) addTrigger(triggerSub *triggerSubscription) error { + respChan := make(chan response) + mqt.reqChan <- request{ + requestType: ADD_TRIGGER, + triggerSub: triggerSub, + respChan: respChan, + } + r := <-respChan + return r.err +} + +func (mqt *MessageQueueTriggerManager) getAllTriggers() *map[string]messageQueueSubscription { + respChan := make(chan response) + mqt.reqChan <- request{ + requestType: GET_ALL_TRIGGERS, + respChan: respChan, + } + r := <-respChan + return r.triggers +} + +func (mqt *MessageQueueTriggerManager) delTrigger(triggerUid string) { + mqt.reqChan <- request{ + requestType: DELETE_TRIGGER, + triggerSub: &triggerSubscription{ + Metadata: fission.Metadata{ + Uid: triggerUid, + }, + }, + } +} + +func (mqt *MessageQueueTriggerManager) syncTriggers() { + for { + // TODO: handle error + newTriggers, err := mqt.controller.MessageQueueTriggerList(mqt.mqCfg.MQType) + if err != nil { + log.Warnf("Failed to sync message queue trigger from controller: %v", err) + } + // sync trigger from controller + newTriggerMap := map[string]fission.MessageQueueTrigger{} + for _, trigger := range newTriggers { + newTriggerMap[trigger.Uid] = trigger + } + currentTriggerSubMap := mqt.getAllTriggers() + + // register new triggers + for key, trigger := range newTriggerMap { + if _, ok := (*currentTriggerSubMap)[key]; ok { + continue + } + sub, err := mqt.messageQueue.subscribe(trigger) + if err != nil { + log.Warnf("Message queue trigger %s created failed: %v", trigger.Name, err) + continue + } + triggerSub := triggerSubscription{ + Metadata: fission.Metadata{ + Name: trigger.Name, + Uid: trigger.Uid, + }, + funcMeta: trigger.Function, + subscription: sub, + } + err = mqt.addTrigger(&triggerSub) + if err != nil { + log.Warnf("Message queue trigger %s created failed: %v", trigger.Name, err) + continue + } + log.Infof("Message queue trigger %s created", trigger.Name) + } + + // remove old triggers + for _, ts := range *currentTriggerSubMap { + triggerSub := ts.(*triggerSubscription) + if _, ok := newTriggerMap[triggerSub.Uid]; ok { + continue + } + if err := mqt.messageQueue.unsubscribe(triggerSub.subscription); err != nil { + log.Warnf("Message queue trigger %s deleted failed: %v", triggerSub.Name, err) + } else { + mqt.delTrigger(triggerSub.Uid) + log.Infof("Message queue trigger %s deleted", triggerSub.Name) + } + } + time.Sleep(3 * time.Second) + } +} + +func IsTopicValid(mqType string, topic string) bool { + switch mqType { + case NATS: + return isTopicValidForNats(topic) + } + return false +} diff --git a/mqtrigger/messageQueue/nats.go b/mqtrigger/messageQueue/nats.go new file mode 100644 index 00000000..1ca0af41 --- /dev/null +++ b/mqtrigger/messageQueue/nats.go @@ -0,0 +1,136 @@ +/* +Copyright 2016 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package messageQueue + +import ( + "bytes" + "errors" + "fmt" + "io/ioutil" + "net/http" + "strings" + + ns "github.com/nats-io/go-nats-streaming" + nsUtil "github.com/nats-io/nats-streaming-server/util" + log "github.com/sirupsen/logrus" + + "github.com/fission/fission" +) + +const ( + natsClusterID = "fissionMQTrigger" + natsProtocol = "nats://" + natsClientID = "fission" + natsQueueGroup = "fission-messageQueueNatsTrigger" +) + +type ( + Nats struct { + nsConn ns.Conn + routerUrl string + } +) + +func makeNatsMessageQueue(routerUrl string, mqCfg MessageQueueConfig) (MessageQueue, error) { + conn, err := ns.Connect(natsClusterID, natsClientID, ns.NatsURL(mqCfg.Url)) + if err != nil { + return nil, err + } + nats := Nats{ + nsConn: conn, + routerUrl: routerUrl, + } + return nats, nil +} + +func (nats Nats) subscribe(trigger fission.MessageQueueTrigger) (messageQueueSubscription, error) { + subj := trigger.Topic + + if !isTopicValidForNats(subj) { + return nil, errors.New(fmt.Sprintf("Not a valid topic: %s", trigger.Topic)) + } + + opts := []ns.SubscriptionOption{ + // Create a durable subscription to nats, so that triggers could retrieve last unack message. + // https://github.com/nats-io/go-nats-streaming#durable-subscriptions + ns.DurableName(trigger.Uid), + + // Nats-streaming server is auto-ack mode by default. Since we want nats-streaming server to + // resend a message if the trigger does not ack it, we need to enable the manual ack mode, so that + // trigger could choose to ack message or simply drop it depend on the response of function pod. + ns.SetManualAckMode(), + } + sub, err := nats.nsConn.Subscribe(subj, msgHandler(&nats, trigger), opts...) + if err != nil { + return nil, err + } + return sub, nil +} + +func (nats Nats) unsubscribe(subscription messageQueueSubscription) error { + return subscription.(ns.Subscription).Close() +} + +func isTopicValidForNats(topic string) bool { + // nats-streaming does not support wildcard channel. + return nsUtil.IsSubjectValid(topic) +} + +func msgHandler(nats *Nats, trigger fission.MessageQueueTrigger) func(*ns.Msg) { + return func(msg *ns.Msg) { + url := nats.routerUrl + "/" + strings.TrimPrefix(fission.UrlForFunction(&trigger.Function), "/") + log.Printf("Making HTTP request to %v", url) + + headers := map[string]string{ + "X-Fission-MQTrigger-Topic": trigger.Topic, + "X-Fission-MQTrigger-RespTopic": trigger.ResponseTopic, + } + + // Create request + req, err := http.NewRequest("POST", url, bytes.NewReader(msg.Data)) + for k, v := range headers { + req.Header.Add(k, v) + } + + // Make the request + resp, err := http.DefaultClient.Do(req) + if err != nil { + log.Warningf("Request failed: %v", url) + return + } + defer resp.Body.Close() + + body, err := ioutil.ReadAll(resp.Body) + if err != nil { + log.Warningf("Request body error: %v", string(body)) + return + } + if resp.StatusCode != 200 { + log.Printf("Request returned failure: %v", resp.StatusCode) + return + } + // trigger acks message only if a request done successfully + err = msg.Ack() + if err != nil { + log.Warningf("Failed to ack message: %v", err) + } + err = nats.nsConn.Publish(trigger.ResponseTopic, body) + if err != nil { + log.Warningf("Failed to publish message to topic %s: %v", trigger.ResponseTopic, err) + } + } +} diff --git a/resource.go b/resource.go index 8f41bf32..ccfbf35f 100644 --- a/resource.go +++ b/resource.go @@ -28,6 +28,10 @@ func (ht HTTPTrigger) Key() string { return ht.Metadata.Name } +func (mqt MessageQueueTrigger) Key() string { + return mqt.Metadata.Name +} + func (tt TimeTrigger) Key() string { return tt.Metadata.Name } diff --git a/test/mqtrigger/main.js b/test/mqtrigger/main.js new file mode 100644 index 00000000..b791d1af --- /dev/null +++ b/test/mqtrigger/main.js @@ -0,0 +1,6 @@ +module.exports = async function(context) { + return { + status: 200, + body: "Hello, World!" + }; +} \ No newline at end of file diff --git a/test/mqtrigger/stan-pub.go b/test/mqtrigger/stan-pub.go new file mode 100644 index 00000000..33c90f7d --- /dev/null +++ b/test/mqtrigger/stan-pub.go @@ -0,0 +1,112 @@ +// This file originally came from official Nats.io GitHub repository. +// You can reach original file with the following link: +// https://github.com/nats-io/go-nats-streaming/tree/master/examples + +// Copyright 2012-2016 Apcera Inc. All rights reserved. +// +build ignore + +package main + +import ( + "flag" + "fmt" + "log" + "os" + "sync" + "time" + + "github.com/nats-io/go-nats-streaming" +) + +var usageStr = ` +Usage: stan-pub [options] + +Options: + -s, --server NATS Streaming server URL(s) + -c, --cluster NATS Streaming cluster name + -id,--clientid NATS Streaming client ID + -a, --async Asynchronous publish mode +` + +// NOTE: Use tls scheme for TLS, e.g. stan-pub -s tls://demo.nats.io:4443 foo hello +func usage() { + fmt.Printf("%s\n", usageStr) + os.Exit(0) +} + +func main() { + var clusterID string + var clientID string + var async bool + var URL string + + flag.StringVar(&URL, "s", stan.DefaultNatsURL, "The nats server URLs (separated by comma)") + flag.StringVar(&URL, "server", stan.DefaultNatsURL, "The nats server URLs (separated by comma)") + flag.StringVar(&clusterID, "c", "test-cluster", "The NATS Streaming cluster ID") + flag.StringVar(&clusterID, "cluster", "test-cluster", "The NATS Streaming cluster ID") + flag.StringVar(&clientID, "id", "stan-pub", "The NATS Streaming client ID to connect with") + flag.StringVar(&clientID, "clientid", "stan-pub", "The NATS Streaming client ID to connect with") + flag.BoolVar(&async, "a", false, "Publish asynchronously") + flag.BoolVar(&async, "async", false, "Publish asynchronously") + + log.SetFlags(0) + flag.Usage = usage + flag.Parse() + + args := flag.Args() + + if len(args) < 1 { + usage() + } + + sc, err := stan.Connect(clusterID, clientID, stan.NatsURL(URL)) + if err != nil { + log.Fatalf("Can't connect: %v.\nMake sure a NATS Streaming Server is running at: %s", err, URL) + } + defer sc.Close() + + subj, msg := args[0], []byte(args[1]) + + ch := make(chan bool) + var glock sync.Mutex + var guid string + acb := func(lguid string, err error) { + glock.Lock() + log.Printf("Received ACK for guid %s\n", lguid) + defer glock.Unlock() + if err != nil { + log.Fatalf("Error in server ack for guid %s: %v\n", lguid, err) + } + if lguid != guid { + log.Fatalf("Expected a matching guid in ack callback, got %s vs %s\n", lguid, guid) + } + ch <- true + } + + if async != true { + err = sc.Publish(subj, msg) + if err != nil { + log.Fatalf("Error during publish: %v\n", err) + } + log.Printf("Published [%s] : '%s'\n", subj, msg) + } else { + glock.Lock() + guid, err = sc.PublishAsync(subj, msg, acb) + if err != nil { + log.Fatalf("Error during async publish: %v\n", err) + } + glock.Unlock() + if guid == "" { + log.Fatal("Expected non-empty guid to be returned.") + } + log.Printf("Published [%s] : '%s' [guid: %s]\n", subj, msg, guid) + + select { + case <-ch: + break + case <-time.After(5 * time.Second): + log.Fatal("timeout") + } + + } +} diff --git a/test/mqtrigger/stan-sub.go b/test/mqtrigger/stan-sub.go new file mode 100644 index 00000000..7a89e8dc --- /dev/null +++ b/test/mqtrigger/stan-sub.go @@ -0,0 +1,133 @@ +// This file originally came from official Nats.io GitHub repository. +// You can reach original file with the following link: +// https://github.com/nats-io/go-nats-streaming/tree/master/examples + +// Copyright 2012-2016 Apcera Inc. All rights reserved. +// +build ignore + +package main + +import ( + "flag" + "log" + "time" + + "github.com/nats-io/go-nats-streaming" + "github.com/nats-io/go-nats-streaming/pb" +) + +var usageStr = ` +Usage: stan-sub [options] + +Options: + -s, --server NATS Streaming server URL(s) + -c, --cluster NATS Streaming cluster name + -id,--clientid NATS Streaming client ID + +Subscription Options: + --qgroup Queue group + --seq Start at seqno + --all Deliver all available messages + --last Deliver starting with last published message + --since Deliver messages in last interval (e.g. 1s, 1hr) + (for more information: https://golang.org/pkg/time/#ParseDuration) + --durable Durable subscriber name + --unsubscribe Unsubscribe the durable on exit +` + +// NOTE: Use tls scheme for TLS, e.g. stan-sub -s tls://demo.nats.io:4443 foo +func usage() { + log.Fatalf(usageStr) +} + +func printMsg(m *stan.Msg) { + log.Printf("[%s]: '%s'", m.Subject, m.Data) +} + +func main() { + var clusterID string + var clientID string + var showTime bool + var startSeq uint64 + var startDelta string + var deliverAll bool + var deliverLast bool + var durable string + var qgroup string + var unsubscribe bool + var URL string + + // defaultID := fmt.Sprintf("client.%s", nuid.Next()) + + flag.StringVar(&URL, "s", stan.DefaultNatsURL, "The nats server URLs (separated by comma)") + flag.StringVar(&URL, "server", stan.DefaultNatsURL, "The nats server URLs (separated by comma)") + flag.StringVar(&clusterID, "c", "test-cluster", "The NATS Streaming cluster ID") + flag.StringVar(&clusterID, "cluster", "test-cluster", "The NATS Streaming cluster ID") + flag.StringVar(&clientID, "id", "", "The NATS Streaming client ID to connect with") + flag.StringVar(&clientID, "clientid", "", "The NATS Streaming client ID to connect with") + flag.BoolVar(&showTime, "t", false, "Display timestamps") + // Subscription options + flag.Uint64Var(&startSeq, "seq", 0, "Start at sequence no.") + flag.BoolVar(&deliverAll, "all", false, "Deliver all") + flag.BoolVar(&deliverLast, "last", false, "Start with last value") + flag.StringVar(&startDelta, "since", "", "Deliver messages since specified time offset") + flag.StringVar(&durable, "durable", "", "Durable subscriber name") + flag.StringVar(&qgroup, "qgroup", "", "Queue group name") + flag.BoolVar(&unsubscribe, "unsubscribe", false, "Unsubscribe the durable on exit") + + log.SetFlags(0) + flag.Usage = usage + flag.Parse() + + args := flag.Args() + + if clientID == "" { + log.Printf("Error: A unique client ID must be specified.") + usage() + } + if len(args) < 1 { + log.Printf("Error: A subject must be specified.") + usage() + } + + sc, err := stan.Connect(clusterID, clientID, stan.NatsURL(URL)) + if err != nil { + log.Fatalf("Can't connect: %v.\nMake sure a NATS Streaming Server is running at: %s", err, URL) + } + // log.Printf("Connected to %s clusterID: [%s] clientID: [%s]\n", URL, clusterID, clientID) + + subj := args[0] + + exit := make(chan struct{}) + mcb := func(msg *stan.Msg) { + printMsg(msg) + exit <- struct{}{} + } + + startOpt := stan.StartAt(pb.StartPosition_NewOnly) + + if startSeq != 0 { + startOpt = stan.StartAtSequence(startSeq) + } else if deliverLast == true { + startOpt = stan.StartWithLastReceived() + } else if deliverAll == true { + log.Print("subscribing with DeliverAllAvailable") + startOpt = stan.DeliverAllAvailable() + } else if startDelta != "" { + ago, err := time.ParseDuration(startDelta) + if err != nil { + sc.Close() + log.Fatal(err) + } + startOpt = stan.StartAtTimeDelta(ago) + } + + sub, err := sc.QueueSubscribe(subj, qgroup, mcb, startOpt, stan.DurableName(durable)) + if err != nil { + sc.Close() + log.Fatal(err) + } + + <-exit + sub.Unsubscribe() +} diff --git a/test/mqtrigger/test.sh b/test/mqtrigger/test.sh new file mode 100755 index 00000000..c196252d --- /dev/null +++ b/test/mqtrigger/test.sh @@ -0,0 +1,45 @@ +#!/bin/bash + +set -e + +clusterID="fissionMQTrigger" +topic="foo.bar" +resptopic="foo.foo" +expectedRespOutput="[foo.foo]: 'Hello, World!'" +FISSIONDIR=$GOPATH"/src/github.com/fission/fission" + +if [[ -z $NATS_STREAMING_URL ]]; then + echo "'NATS_STREAMING_URL' must not be empty. For example: export NATS_STREAMING_URL=nats://192.168.0.1:4222" + exit 1 +fi + +if [[ -z $FISSION_URL ]]; then + echo "'FISSION_URL' must not be empty. For example: export FISSION_URL=http://10.10.10.10" + exit 1 +fi + +cd $FISSIONDIR"/fission/" +go build +mv fission $FISSIONDIR"/test/mqtrigger" +cd $FISSIONDIR"/test/mqtrigger" + +./fission env create --name nodejs --image fission/node-env +./fission fn create --name hello1 --env nodejs --code main.js --method GET +./fission route create --method GET --url /h1 --function hello1 +./fission mqtrigger create --name h1 --function hello1 --mqtype "nats-streaming" --topic "foo.bar" --resptopic "foo.foo" + +# wait until nats trigger is created +sleep 5 + +go run ./stan-pub.go -s $NATS_STREAMING_URL -c $clusterID -id clientPub $topic "" || exit 1 + +response=$(go run ./stan-sub.go --last -s $NATS_STREAMING_URL -c $clusterID -id clientSub $resptopic 2>&1) + +if [[ "$response" != "$expectedRespOutput" ]]; then + echo "$response is not equal to $expectedRespOutput" + exit 1 +fi + +echo "Subscriber received expected response: $response" + +exit 0 diff --git a/types.go b/types.go index 25b320cf..91575c8f 100644 --- a/types.go +++ b/types.go @@ -53,6 +53,14 @@ type ( Function Metadata `json:"function"` } + MessageQueueTrigger struct { + Metadata `json:"metadata"` + Function Metadata `json:"function"` + MessageQueueType string `json:"messageQueueType"` + Topic string `json:"topic"` + ResponseTopic string `json:"respTopic,omitempty"` + } + // Watch is a specification of Kubernetes watch along with a URL to post events to. Watch struct { Metadata `json:"metadata"`