Implement support for Azure storage message queue triggers (#371)

These commits implement support for consuming messages from an Azure storage queue to trigger Fission functions.

* Add stubbed Azure message queue implementation and modify Helm charts.

This commit stubs an implementation for an Azure storage message queue trigger
that will be completed by future commits.

It also modifies the Helm chart to add support for deploying Fission with an
mqtrigger configured for Azure storage queue triggers.

* Add Azure Go SDK to glide.

This commit adds the Azure Go SDK to glide for the upcoming work to support
Azure storage queue triggers.

* Implement Azure message queue trigger.

This commit implements a message queue trigger based on Azure storage queues.

Required message queue trigger manager environment variables:

* AZURE_STORAGE_ACCOUNT_NAME - the Azure storage account to use.
* AZURE_STORAGE_ACCOUNT_KEY - the Azure storage account key.

When creating a message queue trigger, the topic will be the Azure storage
queue to receive messages from.

* Add CA certificates to fission-bundle.

This commit adds the root CA certificates to the fission-bundle image. This
allows Fission to contact third-party APIs that use HTTPS with root CA
signed certificates.

* Add Makefile to build and test.

This commit adds a simple Makefile for building the client and bundle, running
tests, creating the Docker image, and pushing the Docker image.
This commit is contained in:
Peter Huene
2018-02-09 10:38:06 -08:00
committed by Soam Vasani
parent 46895b03c8
commit 4d0e6af5de
15 changed files with 982 additions and 32 deletions
+2
View File
@@ -1,4 +1,6 @@
fission-bundle/fission-bundle
fission/fission
# Pycharm IDE
.idea
environments/php7/vendor/
+46
View File
@@ -0,0 +1,46 @@
# 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.
.PHONY: test
.DEFAULT_GOAL := build
IMAGE ?= fission/fission-bundle
VERSION ?= 0.3.0
ARCH ?= amd64
OS ?= linux
test:
go test -v $(shell go list ./... | grep -v /examples/ | grep -v /environments/)
build: build-bundle build-client
build-client:
go build -o fission/fission ./fission/*.go
build-bundle:
CGO_ENABLED=0 GOOS=$(OS) GOARCH=$(ARCH) go build -o fission-bundle/fission-bundle ./fission-bundle/*.go
build-image:
docker build --rm --tag "$(IMAGE):$(VERSION)" fission-bundle
install:
go install ./fission
image: build-bundle build-image
image-push: image
docker push "$(IMAGE):$(VERSION)"
clean:
@rm -rf fission-bundle/fission-bundle
@rm -rf fission/fission
+11 -11
View File
@@ -39,7 +39,7 @@ The following table lists the configurable parameters of the Fission chart and t
| Parameter | Description | Default |
| ------------------- | ------------------------------------------ | ------------------------ |
| `serviceType` | Type of service to use | `LoadBalancer`. |
| `serviceType` | Type of service to use | `LoadBalancer` |
| `image` | Fission image | `fission/fission-bundle` |
| `imageTag` | Fission image tag | `alpha20170124` |
| `fetcherImage` | Fission fetcher image | `fission/fetcher` |
@@ -49,18 +49,18 @@ The following table lists the configurable parameters of the Fission chart and t
| `functionNamespace` | Namespace for Fission functions | `fission-function` |
| `builderNamespace` | Namespace for Fission environment builders | `fission-builder` |
* Extra configuration for `fission-all`
| Parameter | Description | Default |
| ---------------------- | --------------------------- | -------------------------- |
| `logger.influxdbAdmin` | Log database admin username | `admin`. |
| `logger.fluentdImage` | Logger fluentd image | `fission/fluentd` |
| `fissionUiImage` | Fission ui image | `fission/fission-ui:0.1.0` |
| `nats.authToken` | Nats streaming auth token | `defaultFissionAuthToken` |
| `nats.clusterID` | Nats streaming clusterID | `fissionMQTrigger` |
| Parameter | Description | Default |
| ------------------------------- | --------------------------- | ---------------------------------------------------------- |
| `logger.influxdbAdmin` | Log database admin username | `admin` |
| `logger.fluentdImage` | Logger fluentd image | `fission/fluentd` |
| `fissionUiImage` | Fission ui image | `fission/fission-ui:0.1.0` |
| `messageQueue` | Message queue type | `nats-streaming` |
| `nats.authToken` | Nats streaming auth token | `defaultFissionAuthToken` |
| `nats.clusterID` | Nats streaming clusterID | `fissionMQTrigger` |
| `azureStorageQueue.accountName` | Azure storage account name | None (required if `messageQueue` is `azure-storage-queue`) |
| `azureStorageQueue.key` | Azure storage access key | None (required if `messageQueue` is `azure-storage-queue`) |
Specify each parameter using the `--set key=value[,key=value]` argument to `helm install`. For example,
+13 -1
View File
@@ -444,6 +444,7 @@ spec:
# serviceAccount: fission-svc
---
{{- if eq .Values.messageQueue.type "nats-streaming" }}
apiVersion: extensions/v1beta1
kind: Deployment
metadata:
@@ -469,6 +470,7 @@ spec:
- containerPort: 4222
hostPort: 4222
protocol: TCP
{{- end }}
---
apiVersion: extensions/v1beta1
@@ -492,9 +494,19 @@ spec:
args: ["--mqt"]
env:
- name: MESSAGE_QUEUE_TYPE
value: nats-streaming
value: {{ .Values.messageQueue.type }}
{{- if eq .Values.messageQueue.type "nats-streaming" }}
- name: MESSAGE_QUEUE_URL
value: nats://{{ .Values.nats.authToken }}@nats-streaming:4222
{{- else if eq .Values.messageQueue.type "azure-storage-queue" }}
- name: AZURE_STORAGE_ACCOUNT_NAME
value: {{ required "An Azure storage account name is required." .Values.azureStorageQueue.accountName }}
- name: AZURE_STORAGE_ACCOUNT_KEY
valueFrom:
secretKeyRef:
name: azure-storage-account-key
key: key
{{- end }}
serviceAccount: fission-svc
---
+14 -1
View File
@@ -7,4 +7,17 @@ metadata:
type: Opaque
data:
username: {{ .Values.logger.influxdbAdmin | b64enc | quote }}
password: {{ randAlphaNum 20 | b64enc | quote }}
password: {{ randAlphaNum 20 | b64enc | quote }}
---
{{- if eq .Values.messageQueue.type "azure-storage-queue" }}
apiVersion: v1
kind: Secret
metadata:
name: azure-storage-account-key
labels:
chart: "{{ .Chart.Name }}-{{ .Chart.Version }}"
type: Opaque
data:
key: {{ required "An Azure storage access key is required." .Values.azureStorageQueue.key | b64enc | quote }}
{{- end }}
+2
View File
@@ -36,6 +36,7 @@ spec:
svc: controller
---
{{- if eq .Values.messageQueue.type "nats-streaming" }}
apiVersion: v1
kind: Service
metadata:
@@ -53,6 +54,7 @@ spec:
{{ end }}
selector:
svc: nats-streaming
{{- end }}
---
apiVersion: v1
+9
View File
@@ -44,12 +44,21 @@ logger:
fluentdImage: fission/fluentd
fluentdImageTag: 0.5.0
## Type of Queue you would like to use
## currently supports nats-streaming, azure-storage-queue
messageQueue:
type: nats-streaming
## Message queue trigger config
### NATS Streaming
nats:
authToken: "defaultFissionAuthToken"
clusterID: "fissionMQTrigger"
## Required if messageQueue type is azure-storage-queue
azureStorageQueue:
key: ""
accountName: ""
## Persist data to a persistent volume.
persistence:
+1
View File
@@ -1,2 +1,3 @@
FROM alpine:3.4
RUN apk add --update ca-certificates
ADD fission-bundle /
+1 -1
View File
@@ -108,7 +108,7 @@ func main() {
// Message queue trigger
mqtNameFlag := cli.StringFlag{Name: "name", Usage: "Message queue Trigger name"}
mqtFnNameFlag := cli.StringFlag{Name: "function", Usage: "Function name"}
mqtMQTypeFlag := cli.StringFlag{Name: "mqtype", Usage: "Message queue type, e.g. nats-streaming (optional; uses \"nats-streaming\" if unspecified)"}
mqtMQTypeFlag := cli.StringFlag{Name: "mqtype", Usage: "Message queue type, e.g. nats-streaming, azure-storage-queue (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)"}
mqtMsgContentType := cli.StringFlag{Name: "contenttype, c", Usage: "Content type of messages that publish to the topic (optional; uses \"application/json\" if unspecified)"}
+3 -1
View File
@@ -48,8 +48,10 @@ func mqtCreate(c *cli.Context) error {
mqType = messageQueue.NATS
case messageQueue.NATS:
mqType = messageQueue.NATS
case messageQueue.ASQ:
mqType = messageQueue.ASQ
default:
fatal("Unknown message queue type, currently only \"nats-streaming\" is supported")
fatal("Unknown message queue type, currently only \"nats-streaming, azure-storage-queue \" is supported")
}
// TODO: check topic availability
Generated
+37 -16
View File
@@ -1,24 +1,28 @@
hash: 93c80adbba10750a7fca4536f93880f152527a8f04eedf40a8ea8cbadeb88779
updated: 2017-12-02T21:38:02.116016618-08:00
hash: 03ac7555eb1f5745e82abe42a07828ffcaac3785a77e7f37a596f30db8189e7f
updated: 2018-02-08T22:40:25.586979477-08:00
imports:
- name: cloud.google.com/go
version: 3b1ae45394a234c385be014e9a488f2bb6eef821
subpackages:
- compute/metadata
- internal
- name: github.com/Azure/azure-sdk-for-go
version: f111fc2fa3861c5fdced76cae4c9c71821969577
subpackages:
- storage
- name: github.com/Azure/go-autorest
version: 58f6f26e200fa5dfb40c9cd1c83f3e2c860d779d
version: d4e6b95c12a08b4de2d48b45d5b4d594e5d32fab
subpackages:
- autorest
- autorest/adal
- autorest/azure
- autorest/date
- name: github.com/coreos/etcd
version: 3ac54be402ffe4e6df505814456d4931508aaf21
version: 6a265731e10a5137b991c1aa3a83ecefdd149d50
subpackages:
- client
- name: github.com/davecgh/go-spew
version: 5215b55f46b2b919f50a1df0eaa5886afe4e3b3d
version: 346938d642f2ec3594ed81d874461961cd0faa76
subpackages:
- spew
- name: github.com/dchest/uniuri
@@ -84,9 +88,9 @@ imports:
- name: github.com/gorilla/handlers
version: 90663712d74cb411cbef281bc1e08c19d1a76145
- name: github.com/gorilla/mux
version: 7f08801859139f86dfafd1c296e2cba9a80d292e
version: 53c1911da2b537f792e7cafcb446b05ffe33b996
- name: github.com/graymeta/stow
version: 1c51f76db54d79c9db24cb131011632eb586196c
version: abb68c488872b06c5453865fa59f4818b4ea13a4
subpackages:
- local
- name: github.com/hashicorp/golang-lru
@@ -111,10 +115,12 @@ imports:
- buffer
- jlexer
- jwriter
- name: github.com/marstr/guid
version: 8bdf7d1a087ccc975cf37dd6507da50698fd19ca
- name: github.com/mholt/archiver
version: 26cf5bb32d07aa4e8d0de15f56ce516f4641d7df
- name: github.com/nats-io/go-nats
version: 6f06d34c19fc8607adf14668b5bb2c9465d78348
version: 78ec4b93936d7a00e59b7f3939d8116955916d1b
subpackages:
- encoders/builtin
- util
@@ -123,15 +129,16 @@ imports:
subpackages:
- pb
- name: github.com/nats-io/nats-streaming-server
version: f9638df9ef7ee548313eeb5a9dda9c888dc769a0
version: 6fdcdfbb2589e68692a68fd5b64b1b7e7c54bf05
subpackages:
- spb
- util
- name: github.com/nats-io/nuid
version: 289cccf02c178dc782430d534e3c1f5b72af807f
- name: github.com/nwaples/rardecode
version: e06696f847aeda6f39a8f0b7cdff193b7690aef6
- name: github.com/pierrec/lz4
version: 08c27939df1bd95e881e2c2367a749964ad1fceb
version: ed8d4cc3b461464e69798080a0092bd028910298
- name: github.com/pierrec/xxHash
version: a0006b13c722f7f12368c00a3d3c2ae8a999a0c6
subpackages:
@@ -143,13 +150,19 @@ imports:
- name: github.com/PuerkitoBio/urlesc
version: 5bd2802263f21d8788851d5305584c82a5c75d7e
- name: github.com/robfig/cron
version: 736158dc09e10f1911ca3a1e1b01f11b566ce5db
version: 2315d5715e36303a941d907f038da7f7c44c773b
- name: github.com/satori/go.uuid
version: 879c5887cd475cd7864858769793b2ceb0d44feb
version: f58768cc1a7a7e77a3bd49e98cdd21419399b6a3
- name: github.com/sirupsen/logrus
version: 68cec9f21fbf3ea8d8f98c044bc6ce05f17b267a
- name: github.com/spf13/pflag
version: 9ff6c6923cfffbcd502984b8e0c80539a94968b7
- name: github.com/stretchr/testify
version: 12b6f73e6084dad08a7c6e575284b177ecafbc71
subpackages:
- assert
- mock
- require
- name: github.com/ugorji/go
version: ded73eae5db7e7a0ef6f55aace87a2873c5d2b74
subpackages:
@@ -200,8 +213,10 @@ imports:
- unicode/norm
- width
- name: google.golang.org/appengine
version: 12d5545dc1cfa6047a286d5e853841b6471f4c19
version: 9d8544a6b2c7df9cff240fcf92d7b2f59bc13416
repo: https://github.com/golang/appengine
subpackages:
- cloudsql
- internal
- internal/app_identity
- internal/base
@@ -218,7 +233,7 @@ imports:
- name: k8s.io/api
version: 4b8fc5be9b77d91bbb6525d18591c43699a2b4e5
- name: k8s.io/apiextensions-apiserver
version: fcd622fe88a4a6efcb5aea9e94ee87324ac1b036
version: 0965a40c0530e110459750a7fcb9cfa4910fb944
subpackages:
- pkg/apis/apiextensions
- pkg/apis/apiextensions/v1beta1
@@ -226,7 +241,7 @@ imports:
- pkg/client/clientset/clientset/scheme
- pkg/client/clientset/clientset/typed/apiextensions/v1beta1
- name: k8s.io/apimachinery
version: 8ab5f3d8a330c2e9baaf84e39042db8d49034ae2
version: 80184f5f100c67bfcda6b8a848e0243fe4bc2772
subpackages:
- pkg/api/equality
- pkg/api/errors
@@ -358,4 +373,10 @@ imports:
- util/homedir
- util/integer
- util/jsonpath
testImports: []
testImports:
- name: github.com/pmezard/go-difflib
version: d8ed2627bdf02c080bf22230dbb337003b7aba2d
subpackages:
- difflib
- name: github.com/stretchr/objx
version: 8a3f7159479fbc75b30357fbc48f380b7320f08e
+8 -1
View File
@@ -50,4 +50,11 @@ import:
- package: github.com/mholt/archiver
- package: github.com/pkg/errors
- package: github.com/fsnotify/fsnotify
- package: github.com/Azure/azure-sdk-for-go
version: ~12.4.0-beta
- package: github.com/Azure/go-autorest
version: ~9.9.0
- package: github.com/stretchr/testify
version: ~1.2.1
- package: github.com/davecgh/go-spew
version: ~1.1.0
+376
View File
@@ -0,0 +1,376 @@
/*
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 messageQueue
import (
"bytes"
"encoding/base64"
"errors"
"fmt"
"io/ioutil"
"net/http"
"os"
"strconv"
"strings"
"sync"
"time"
"github.com/fission/fission"
"github.com/fission/fission/crd"
log "github.com/sirupsen/logrus"
"github.com/Azure/azure-sdk-for-go/storage"
)
// TODO: some of these constants should probably be environment variables
const (
// AzureQueuePollingInterval is the polling interval (default is 1 minute).
AzureQueuePollingInterval = time.Minute
// AzureQueueRetryLimit is the limit for attempts to retry invoking a function.
AzureQueueRetryLimit = 3
// AzureMessageFetchCount is the number of messages to fetch at a time.
AzureMessageFetchCount = 10
// AzureMessageVisibilityTimeout is the visibility timeout for dequeued messages.
AzureMessageVisibilityTimeout = time.Minute
// AzurePoisonQueueSuffix is the suffix used for posion queues.
AzurePoisonQueueSuffix = "-poison"
// AzureFunctionInvocationTimeout is the amount of time to wait for a triggered function to execute.
AzureFunctionInvocationTimeout = 10 * time.Minute
)
// AzureStorageConnection represents an Azure storage connection.
type AzureStorageConnection struct {
routerURL string
service AzureQueueService
httpClient AzureHTTPClient
}
// AzureQueueSubscription represents an Azure storage message queue subscription.
type AzureQueueSubscription struct {
queue AzureQueue
queueName string
outputQueueName string
functionURL string
contentType string
unsubscribe chan bool
done chan bool
}
// AzureQueueService is the interface that abstracts the Azure storage service.
// This exists to enable unit testing.
type AzureQueueService interface {
GetQueue(name string) AzureQueue
}
// AzureQueue is the interface that abstracts Azure storage queues.
// This exists to enable unit testing.
type AzureQueue interface {
Create(options *storage.QueueServiceOptions) error
NewMessage(text string) AzureMessage
GetMessages(options *storage.GetMessagesOptions) ([]AzureMessage, error)
}
// AzureMessage is the interface that abstracts Azure storage messages.
// This exists to enable unit testing.
type AzureMessage interface {
Bytes() []byte
Put(options *storage.PutMessageOptions) error
Delete(options *storage.QueueServiceOptions) error
}
// AzureHTTPClient is the interface that abstract HTTP requests made by the trigger.
// This exists to enable unit testing.
type AzureHTTPClient interface {
Do(req *http.Request) (*http.Response, error)
}
type azureQueueService struct {
service storage.QueueServiceClient
}
func (qs azureQueueService) GetQueue(name string) AzureQueue {
return azureQueue{
ref: qs.service.GetQueueReference(name),
}
}
type azureQueue struct {
ref *storage.Queue
}
func (qr azureQueue) Create(options *storage.QueueServiceOptions) error {
exists, err := qr.ref.Exists()
if err != nil {
return err
}
if exists {
return nil
}
return qr.ref.Create(options)
}
func (qr azureQueue) NewMessage(text string) AzureMessage {
return azureMessage{
ref: qr.ref.GetMessageReference(text),
bytes: []byte(text),
}
}
func (qr azureQueue) GetMessages(options *storage.GetMessagesOptions) ([]AzureMessage, error) {
msgs, err := qr.ref.GetMessages(options)
if err != nil {
return nil, err
}
messages := make([]AzureMessage, len(msgs))
for i := range msgs {
bytes, err := base64.StdEncoding.DecodeString(msgs[i].Text)
if err != nil {
return nil, err
}
messages[i] = azureMessage{
ref: &msgs[i],
bytes: bytes,
}
}
return messages, nil
}
type azureMessage struct {
ref *storage.Message
bytes []byte
}
func (m azureMessage) Bytes() []byte {
return m.bytes
}
func (m azureMessage) Put(options *storage.PutMessageOptions) error {
return m.ref.Put(options)
}
func (m azureMessage) Delete(options *storage.QueueServiceOptions) error {
return m.ref.Delete(options)
}
func newAzureQueueService(client storage.Client) AzureQueueService {
return azureQueueService{
service: client.GetQueueService(),
}
}
func newAzureStorageConnection(routerURL string, config MessageQueueConfig) (MessageQueue, error) {
account := os.Getenv("AZURE_STORAGE_ACCOUNT_NAME")
if len(account) == 0 {
return nil, errors.New("Required environment variable 'AZURE_STORAGE_ACCOUNT_NAME' is not set")
}
key := os.Getenv("AZURE_STORAGE_ACCOUNT_KEY")
if len(key) == 0 {
return nil, errors.New("Required environment variable 'AZURE_STORAGE_ACCOUNT_KEY' is not set")
}
log.Infof("Creating Azure storage connection to storage account '%s'.", account)
client, err := storage.NewBasicClient(account, key)
if err != nil {
return nil, fmt.Errorf("Failed to Azure create storage client: %v", err)
}
return &AzureStorageConnection{
routerURL: routerURL,
service: newAzureQueueService(client),
httpClient: &http.Client{
Timeout: AzureFunctionInvocationTimeout,
},
}, nil
}
func (asc AzureStorageConnection) subscribe(trigger *crd.MessageQueueTrigger) (messageQueueSubscription, error) {
log.Infof("Subscribing to Azure storage queue '%s'.", trigger.Spec.Topic)
if trigger.Spec.FunctionReference.Type != fission.FunctionReferenceTypeFunctionName {
return nil, fmt.Errorf("Unsupported function reference type (%v) for trigger %v", trigger.Spec.FunctionReference.Type, trigger.Metadata.Name)
}
subscription := &AzureQueueSubscription{
queue: asc.service.GetQueue(trigger.Spec.Topic),
queueName: trigger.Spec.Topic,
outputQueueName: trigger.Spec.ResponseTopic,
functionURL: asc.routerURL + "/" + strings.TrimPrefix(fission.UrlForFunction(trigger.Spec.FunctionReference.Name), "/"),
contentType: trigger.Spec.ContentType,
unsubscribe: make(chan bool),
done: make(chan bool),
}
go runAzureQueueSubscription(asc, subscription)
return subscription, nil
}
func (asc AzureStorageConnection) unsubscribe(subscription messageQueueSubscription) error {
sub := subscription.(*AzureQueueSubscription)
log.Infof("Unsubscribing from Azure storage queue '%s'.", sub.queueName)
// Let the worker know we've unsubscribed
sub.unsubscribe <- true
// Wait until the subscription is done
<-sub.done
return nil
}
func runAzureQueueSubscription(conn AzureStorageConnection, sub *AzureQueueSubscription) {
var wg sync.WaitGroup
// Process the queue before waiting
pollAzureQueueSubscription(conn, sub, &wg)
timer := time.NewTimer(AzureQueuePollingInterval)
for {
log.Infof("Waiting for %v before polling Azure storage queue '%s'.", AzureQueuePollingInterval, sub.queueName)
select {
case <-sub.unsubscribe:
timer.Stop()
wg.Wait()
sub.done <- true
return
case <-timer.C:
pollAzureQueueSubscription(conn, sub, &wg)
timer.Reset(AzureQueuePollingInterval)
continue
}
}
}
func pollAzureQueueSubscription(conn AzureStorageConnection, sub *AzureQueueSubscription, wg *sync.WaitGroup) {
log.Infof("Polling messages for Azure storage queue '%s'.", sub.queueName)
err := sub.queue.Create(nil)
if err != nil {
log.Errorf("Failed to create message queue '%s': %v", sub.queueName, err)
return
}
for {
err := sub.queue.Create(nil)
if err != nil {
log.Errorf("Failed to create message queue '%s': %v", sub.queueName, err)
return
}
messages, err := sub.queue.GetMessages(&storage.GetMessagesOptions{
NumOfMessages: AzureMessageFetchCount,
VisibilityTimeout: int(AzureMessageVisibilityTimeout / time.Second),
})
if err != nil {
log.Errorf("Failed to retrieve messages from Azure storage queue '%s': %v", sub.queueName, err)
break
}
if len(messages) == 0 {
break
}
wg.Add(len(messages))
for _, msg := range messages {
go func(conn AzureStorageConnection, sub *AzureQueueSubscription, msg AzureMessage) {
defer wg.Done()
invokeTriggeredFunction(conn, sub, msg)
}(conn, sub, msg)
}
}
}
func invokeTriggeredFunction(conn AzureStorageConnection, sub *AzureQueueSubscription, message AzureMessage) {
defer message.Delete(nil)
log.Printf("Making HTTP request to %s.", sub.functionURL)
for i := 0; i <= AzureQueueRetryLimit; i++ {
if i > 0 {
log.Infof("Retry #%d for request to %s.", i, sub.functionURL)
}
request, err := http.NewRequest("POST", sub.functionURL, bytes.NewReader(message.Bytes()))
if err != nil {
log.Errorf("Failed to create HTTP request to %s: %v", sub.functionURL, err)
continue
}
request.Header.Add("X-Fission-MQTrigger-Topic", sub.queueName)
if len(sub.outputQueueName) > 0 {
request.Header.Add("X-Fission-MQTrigger-RespTopic", sub.outputQueueName)
}
if i > 0 {
request.Header.Add("X-Fission-MQTrigger-RetryCount", strconv.Itoa(i))
}
request.Header.Add("Content-Type", sub.contentType)
response, err := conn.httpClient.Do(request)
if err != nil {
log.Errorf("Request to %s failed: %v", sub.functionURL, err)
continue
}
defer response.Body.Close()
body, err := ioutil.ReadAll(response.Body)
if err != nil {
log.Errorf("Failed to read response body from %s: %v.", sub.functionURL, err)
continue
}
if response.StatusCode < 200 || response.StatusCode >= 300 {
log.Printf("Request to %s returned failure: %s (%d).", sub.functionURL, string(body), response.StatusCode)
continue
}
if len(sub.outputQueueName) > 0 {
outputQueue := conn.service.GetQueue(sub.outputQueueName)
err = outputQueue.Create(nil)
if err != nil {
log.Errorf("Failed to create output queue '%s': %v.", sub.outputQueueName, err)
return
}
outputMessage := outputQueue.NewMessage(string(body))
err = outputMessage.Put(nil)
if err != nil {
log.Errorf("Failed to post response body from %s to output queue '%s': %v.", sub.functionURL, sub.outputQueueName, err)
return
}
}
// Function invocation was successful
return
}
log.Errorf("Request to %s failed after %d retries; moving message to poison queue.", sub.functionURL, AzureQueueRetryLimit)
poisonQueueName := sub.queueName + AzurePoisonQueueSuffix
poisonQueue := conn.service.GetQueue(poisonQueueName)
err := poisonQueue.Create(nil)
if err != nil {
log.Errorf("Failed to create poison queue '%s': %v", poisonQueueName, err)
return
}
poisonMessage := poisonQueue.NewMessage(string(message.Bytes()))
err = poisonMessage.Put(nil)
if err != nil {
log.Errorf("Failed to post response body from %s to output queue '%s': %v", sub.functionURL, poisonQueueName, err)
return
}
}
+449
View File
@@ -0,0 +1,449 @@
/*
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 messageQueue
import (
"fmt"
"io/ioutil"
"net/http"
"os"
"strings"
"testing"
"time"
"github.com/fission/fission"
"github.com/fission/fission/crd"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"github.com/Azure/azure-sdk-for-go/storage"
"github.com/stretchr/testify/mock"
"github.com/stretchr/testify/require"
)
const (
DummyRouterURL = "http://localhost"
)
type azureQueueServiceMock struct {
mock.Mock
}
func (m *azureQueueServiceMock) GetQueue(name string) AzureQueue {
args := m.Called(name)
return args.Get(0).(AzureQueue)
}
type azureQueueMock struct {
mock.Mock
}
func (m *azureQueueMock) Create(options *storage.QueueServiceOptions) error {
args := m.Called(options)
return args.Error(0)
}
func (m *azureQueueMock) NewMessage(text string) AzureMessage {
args := m.Called(text)
return args.Get(0).(AzureMessage)
}
func (m *azureQueueMock) GetMessages(options *storage.GetMessagesOptions) ([]AzureMessage, error) {
args := m.Called(options)
return args.Get(0).([]AzureMessage), args.Error(1)
}
type azureMessageMock struct {
mock.Mock
}
func (m *azureMessageMock) Bytes() []byte {
args := m.Called()
return args.Get(0).([]byte)
}
func (m *azureMessageMock) Put(options *storage.PutMessageOptions) error {
args := m.Called(options)
return args.Error(0)
}
func (m *azureMessageMock) Delete(options *storage.QueueServiceOptions) error {
args := m.Called(options)
return args.Error(0)
}
type azureHTTPClientMock struct {
mock.Mock
bodyHandler func(res *http.Response)
}
func (m *azureHTTPClientMock) Do(req *http.Request) (*http.Response, error) {
args := m.Called(req)
res := args.Get(0).(*http.Response)
err := args.Error(1)
if res != nil && m.bodyHandler != nil {
m.bodyHandler(res)
}
return res, err
}
func TestNewStorageConnectionMissingAccountName(t *testing.T) {
connection, err := newAzureStorageConnection(DummyRouterURL, MessageQueueConfig{
MQType: ASQ,
Url: "",
})
require.Nil(t, connection)
require.Error(t, err, "Required environment variable 'AZURE_STORAGE_ACCOUNT_NAME' is not set")
}
func TestNewStorageConnectionMissingAccessKey(t *testing.T) {
_ = os.Setenv("AZURE_STORAGE_ACCOUNT_NAME", "accountname")
connection, err := newAzureStorageConnection(DummyRouterURL, MessageQueueConfig{
MQType: ASQ,
Url: "",
})
_ = os.Unsetenv("AZURE_STORAGE_ACCOUNT_NAME")
require.Nil(t, connection)
require.Error(t, err, "Required environment variable 'AZURE_STORAGE_ACCOUNT_KEY' is not set")
}
func TestNewStorageConnection(t *testing.T) {
_ = os.Setenv("AZURE_STORAGE_ACCOUNT_NAME", "accountname")
_ = os.Setenv("AZURE_STORAGE_ACCOUNT_KEY", "bm90IGEga2V5")
connection, err := newAzureStorageConnection(DummyRouterURL, MessageQueueConfig{
MQType: "azure-storage-queue",
Url: "",
})
_ = os.Unsetenv("AZURE_STORAGE_ACCOUNT_NAME")
_ = os.Unsetenv("AZURE_STORAGE_ACCOUNT_KEY")
require.NoError(t, err)
require.IsType(t, &AzureStorageConnection{}, connection)
p := connection.(*AzureStorageConnection)
require.Equal(t, DummyRouterURL, p.routerURL)
require.NotNil(t, p.service)
}
func TestAzureStorageQueueSingleMessage(t *testing.T) {
runAzureStorageQueueTest(t, 1, false)
}
func TestAzureStorageQueueMultipleMessages(t *testing.T) {
runAzureStorageQueueTest(t, 10, false)
}
func TestAzureStorageQueueSingleOutputMessage(t *testing.T) {
runAzureStorageQueueTest(t, 1, true)
}
func TestAzureStorageQueueMultipleOutputMessages(t *testing.T) {
runAzureStorageQueueTest(t, 10, true)
}
func TestAzureStorageQueuePoisonMessage(t *testing.T) {
const (
TriggerName = "queuetrigger"
QueueName = "inputqueue"
MessageBody = "input"
FunctionName = "badfunc"
ContentType = "text/plain"
)
// Mock a HTTP client that returns different failures
httpClient := new(azureHTTPClientMock)
httpClient.On(
"Do",
mock.MatchedBy(httpRequestMatcher(t, QueueName, "", "", ContentType, FunctionName, MessageBody)),
).Return(
&http.Response{
StatusCode: 500,
Body: ioutil.NopCloser(strings.NewReader("server error")),
},
nil,
).Once()
httpClient.On(
"Do",
mock.MatchedBy(httpRequestMatcher(t, QueueName, "", "1", ContentType, FunctionName, MessageBody)),
).Return(
&http.Response{
StatusCode: 404,
Body: ioutil.NopCloser(strings.NewReader("not found")),
},
nil,
).Once()
httpClient.On(
"Do",
mock.MatchedBy(httpRequestMatcher(t, QueueName, "", "2", ContentType, FunctionName, MessageBody)),
).Return(
&http.Response{
StatusCode: 400,
Body: ioutil.NopCloser(strings.NewReader("bad request")),
},
nil,
).Once()
httpClient.On(
"Do",
mock.MatchedBy(httpRequestMatcher(t, QueueName, "", "3", ContentType, FunctionName, MessageBody)),
).Return(
&http.Response{
StatusCode: 403,
Body: ioutil.NopCloser(strings.NewReader("not authorized")),
},
nil,
).Once()
// Mock a queue message with "input" as the message body
message := new(azureMessageMock)
message.On("Bytes").Return([]byte(MessageBody))
message.On(
"Delete",
mock.MatchedBy(
func(options *storage.QueueServiceOptions) bool {
return options == nil
},
),
).Return(nil)
// Mock a queue that performs a no-op create, returns a "poison" message, and then returns no more messages
queue := new(azureQueueMock)
queue.On(
"Create",
mock.MatchedBy(
func(options *storage.QueueServiceOptions) bool {
return options == nil
},
),
).Return(nil)
queue.On(
"GetMessages",
mock.MatchedBy(
func(options *storage.GetMessagesOptions) bool {
return options.NumOfMessages == AzureMessageFetchCount &&
options.VisibilityTimeout == int(AzureMessageVisibilityTimeout/time.Second)
},
),
).Return([]AzureMessage{message}, nil).Once()
queue.On(
"GetMessages",
mock.MatchedBy(
func(options *storage.GetMessagesOptions) bool {
return options.NumOfMessages == AzureMessageFetchCount &&
options.VisibilityTimeout == int(AzureMessageVisibilityTimeout/time.Second)
},
),
).Return([]AzureMessage{}, nil)
// Mock a poison queue message that performs a no-op Put
poisonMessage := new(azureMessageMock)
poisonMessage.On(
"Put",
mock.MatchedBy(
func(options *storage.PutMessageOptions) bool {
return options == nil
},
),
).Return(nil)
// Mock a poison queue that performs a no-op create and creates a new message
poisonQueue := new(azureQueueMock)
poisonQueue.On(
"Create",
mock.MatchedBy(
func(options *storage.QueueServiceOptions) bool {
return options == nil
},
),
).Return(nil)
poisonQueue.On("NewMessage", MessageBody).Return(poisonMessage).Once()
// Mock the queue service to return the input queue
service := new(azureQueueServiceMock)
service.On("GetQueue", QueueName).Return(queue).Once()
service.On("GetQueue", QueueName+AzurePoisonQueueSuffix).Return(poisonQueue).Once()
// Create the storage connection and subscribe to the trigger
connection := AzureStorageConnection{
routerURL: DummyRouterURL,
service: service,
httpClient: httpClient,
}
subscription, err := connection.subscribe(&crd.MessageQueueTrigger{
Metadata: metav1.ObjectMeta{
Name: TriggerName,
},
Spec: fission.MessageQueueTriggerSpec{
FunctionReference: fission.FunctionReference{
Type: fission.FunctionReferenceTypeFunctionName,
Name: FunctionName,
},
MessageQueueType: ASQ,
Topic: QueueName,
ContentType: ContentType,
},
})
require.NoError(t, err)
require.NotNil(t, subscription)
connection.unsubscribe(subscription)
mock.AssertExpectationsForObjects(t, httpClient, message, poisonMessage, queue, poisonQueue, service)
}
func httpRequestMatcher(t *testing.T, queue string, responseQueue string, retry string, contentType string, functionName string, body string) func(*http.Request) bool {
expectedURL := fmt.Sprintf("%s/fission-function/%s", DummyRouterURL, functionName)
return func(req *http.Request) bool {
requestBody, err := ioutil.ReadAll(req.Body)
require.NoError(t, err)
req.Body = ioutil.NopCloser(strings.NewReader(string(requestBody)))
return queue == req.Header.Get("X-Fission-MQTrigger-Topic") &&
responseQueue == req.Header.Get("X-Fission-MQTrigger-RespTopic") &&
retry == req.Header.Get("X-Fission-MQTrigger-RetryCount") &&
contentType == req.Header.Get("Content-Type") &&
req.URL.String() == expectedURL &&
string(requestBody) == body
}
}
func runAzureStorageQueueTest(t *testing.T, count int, output bool) {
const (
TriggerName = "queuetrigger"
QueueName = "inputqueue"
OutputQueueName = "outputqueue"
MessageBody = "input"
FunctionName = "testfunc"
FunctionResponse = "output"
ContentType = "text/plain"
)
responseTopic := ""
if output {
responseTopic = OutputQueueName
}
// Mock a HTTP client that returns 200 with "output" for the body
httpClient := new(azureHTTPClientMock)
httpClient.bodyHandler = func(res *http.Response) {
res.Body = ioutil.NopCloser(strings.NewReader(FunctionResponse))
}
httpClient.On(
"Do",
mock.MatchedBy(httpRequestMatcher(t, QueueName, responseTopic, "", ContentType, FunctionName, MessageBody)),
).Return(&http.Response{StatusCode: 200}, nil).Times(count)
// Mock a queue message with "input" as the message body
message := new(azureMessageMock)
message.On("Bytes").Return([]byte(MessageBody)).Times(count)
message.On(
"Delete",
mock.MatchedBy(
func(options *storage.QueueServiceOptions) bool {
return options == nil
},
),
).Return(nil).Times(count)
// Mock a queue that performs a no-op create, returns a message the specified number of times, and then returns no more messages
queue := new(azureQueueMock)
queue.On(
"Create",
mock.MatchedBy(
func(options *storage.QueueServiceOptions) bool {
return options == nil
},
),
).Return(nil)
queue.On(
"GetMessages",
mock.MatchedBy(
func(options *storage.GetMessagesOptions) bool {
return options.NumOfMessages == AzureMessageFetchCount &&
options.VisibilityTimeout == int(AzureMessageVisibilityTimeout/time.Second)
},
),
).Return([]AzureMessage{message}, nil).Times(count)
queue.On(
"GetMessages",
mock.MatchedBy(
func(options *storage.GetMessagesOptions) bool {
return options.NumOfMessages == AzureMessageFetchCount &&
options.VisibilityTimeout == int(AzureMessageVisibilityTimeout/time.Second)
},
),
).Return([]AzureMessage{}, nil)
// Mock the output queue if needed
outputMessage := new(azureMessageMock)
outputQueue := new(azureQueueMock)
if output {
outputMessage.On(
"Put",
mock.MatchedBy(
func(options *storage.PutMessageOptions) bool {
return options == nil
},
),
).Return(nil).Times(count)
outputQueue.On(
"Create",
mock.MatchedBy(
func(options *storage.QueueServiceOptions) bool {
return options == nil
},
),
).Return(nil).Times(count)
outputQueue.On("NewMessage", FunctionResponse).Return(outputMessage).Times(count)
}
// Mock the queue service to return the input queue
service := new(azureQueueServiceMock)
service.On("GetQueue", QueueName).Return(queue).Once()
if output {
service.On("GetQueue", OutputQueueName).Return(outputQueue).Times(count)
}
// Create the storage connection and subscribe to the trigger
connection := AzureStorageConnection{
routerURL: DummyRouterURL,
service: service,
httpClient: httpClient,
}
subscription, err := connection.subscribe(&crd.MessageQueueTrigger{
Metadata: metav1.ObjectMeta{
Name: TriggerName,
},
Spec: fission.MessageQueueTriggerSpec{
FunctionReference: fission.FunctionReference{
Type: fission.FunctionReferenceTypeFunctionName,
Name: FunctionName,
},
MessageQueueType: ASQ,
Topic: QueueName,
ResponseTopic: responseTopic,
ContentType: ContentType,
},
})
require.NoError(t, err)
require.NotNil(t, subscription)
connection.unsubscribe(subscription)
mock.AssertExpectationsForObjects(t, httpClient, message, outputMessage, queue, outputQueue, service)
}
+10
View File
@@ -18,6 +18,7 @@ package messageQueue
import (
"errors"
"regexp"
"time"
log "github.com/sirupsen/logrus"
@@ -28,6 +29,7 @@ import (
const (
NATS string = "nats-streaming"
ASQ string = "azure-storage-queue"
)
const (
@@ -36,6 +38,10 @@ const (
GET_ALL_TRIGGERS
)
var (
validAzureQueueName = regexp.MustCompile("^[a-z0-9][a-z0-9\\-]*[a-z0-9]$")
)
type (
messageQueueSubscription interface{}
@@ -87,6 +93,8 @@ func MakeMessageQueueTriggerManager(fissionClient *crd.FissionClient, routerUrl
switch mqConfig.MQType {
case NATS:
messageQueue, err = makeNatsMessageQueue(routerUrl, mqConfig)
case ASQ:
messageQueue, err = newAzureStorageConnection(routerUrl, mqConfig)
default:
err = errors.New("No matched message queue type found")
}
@@ -222,6 +230,8 @@ func IsTopicValid(mqType string, topic string) bool {
switch mqType {
case NATS:
return isTopicValidForNats(topic)
case ASQ:
return len(topic) >= 3 && len(topic) <= 63 && validAzureQueueName.MatchString(topic)
}
return false
}