From 4e91579ef2b5ab301410636f71b2d0707e24873a Mon Sep 17 00:00:00 2001 From: Sanket Sudake Date: Wed, 6 Apr 2022 13:41:38 +0530 Subject: [PATCH] Remove deprecated Fission Nats connector (#2403) We are removing Fission deprecated Nats connector and planning to adopt Keda going forward to have better delegated functionality and more rich support. Signed-off-by: Sanket Sudake --- .github/workflows/push_pr.yaml | 1 - .../mqt-fission-nats/deployment.yaml | 120 --------- .../templates/mqt-fission-nats/svc.yaml | 19 -- charts/fission-all/values.yaml | 45 ---- cmd/fission-bundle/mqtrigger/mqtrigger.go | 1 - cmd/fission-cli/app/app.go | 1 - go.mod | 9 +- go.sum | 59 ----- pkg/apis/core/v1/const.go | 1 - pkg/fission-cli/cmd/support/dump.go | 10 +- pkg/fission-cli/cmd/support/resources/crd.go | 2 +- pkg/fission-cli/flag/flag.go | 2 +- pkg/mqtrigger/messageQueue/nats/nats.go | 245 ------------------ skaffold.yaml | 1 - test/kind_CI.sh | 4 +- test/test_utils.sh | 4 - test/tests/mqtrigger/nats/main.js | 6 - test/tests/mqtrigger/nats/main_error.js | 6 - test/tests/mqtrigger/nats/stan-pub/main.go | 143 ---------- test/tests/mqtrigger/nats/stan-sub/main.go | 188 -------------- test/tests/mqtrigger/nats/test_mqtrigger.sh | 72 ----- .../mqtrigger/nats/test_mqtrigger_error.sh | 75 ------ tools/port-forward-nats.sh | 23 -- 23 files changed, 9 insertions(+), 1028 deletions(-) delete mode 100644 charts/fission-all/templates/mqt-fission-nats/deployment.yaml delete mode 100644 charts/fission-all/templates/mqt-fission-nats/svc.yaml delete mode 100644 pkg/mqtrigger/messageQueue/nats/nats.go delete mode 100644 test/tests/mqtrigger/nats/main.js delete mode 100644 test/tests/mqtrigger/nats/main_error.js delete mode 100644 test/tests/mqtrigger/nats/stan-pub/main.go delete mode 100644 test/tests/mqtrigger/nats/stan-sub/main.go delete mode 100755 test/tests/mqtrigger/nats/test_mqtrigger.sh delete mode 100755 test/tests/mqtrigger/nats/test_mqtrigger_error.sh delete mode 100755 tools/port-forward-nats.sh diff --git a/.github/workflows/push_pr.yaml b/.github/workflows/push_pr.yaml index 6ecc0aa2..25c9d94a 100644 --- a/.github/workflows/push_pr.yaml +++ b/.github/workflows/push_pr.yaml @@ -108,7 +108,6 @@ jobs: run: | kubectl port-forward svc/router 8888:80 -nfission & kubectl port-forward svc/controller 8889:80 -nfission & - kubectl port-forward svc/nats-streaming 8890:4222 -nfission & - name: Get fission version run: | diff --git a/charts/fission-all/templates/mqt-fission-nats/deployment.yaml b/charts/fission-all/templates/mqt-fission-nats/deployment.yaml deleted file mode 100644 index e16596e5..00000000 --- a/charts/fission-all/templates/mqt-fission-nats/deployment.yaml +++ /dev/null @@ -1,120 +0,0 @@ -{{- if .Values.nats.enabled }} -{{- if not .Values.nats.external }} -apiVersion: v1 -kind: ServiceAccount -metadata: - name: fission-nats-streaming - namespace: {{ .Release.Namespace }} ---- -apiVersion: apps/v1 -kind: Deployment -metadata: - labels: - svc: nats-streaming - name: nats-streaming -spec: - replicas: 1 - selector: - matchLabels: - svc: nats-streaming - template: - metadata: - labels: - svc: nats-streaming - spec: - serviceAccount: fission-nats-streaming - containers: - - name: nats-streaming - image: "{{ .Values.nats.streamingserver.image }}:{{ .Values.nats.streamingserver.tag }}" - imagePullPolicy: {{ .Values.pullPolicy }} - args: [ - "--cluster_id", "{{ .Values.nats.clusterID }}", - "--auth", "{{ .Values.nats.authToken }}", - "--max_channels", "0", - "--http_port", "4223" - ] - ports: - - containerPort: 4222 - protocol: TCP - - containerPort: 4223 - protocol: TCP - readinessProbe: - httpGet: - path: "/streaming/serverz" - port: 4223 - initialDelaySeconds: 30 - periodSeconds: 1 - failureThreshold: 30 - livenessProbe: - httpGet: - path: "/streaming/serverz" - port: 4223 - initialDelaySeconds: 30 - periodSeconds: 5 - {{- if .Values.terminationMessagePath }} - terminationMessagePath: {{ .Values.terminationMessagePath }} - {{- end }} - {{- if .Values.terminationMessagePolicy }} - terminationMessagePolicy: {{ .Values.terminationMessagePolicy }} - {{- end }} -{{- if .Values.extraCoreComponentPodConfig }} -{{ toYaml .Values.extraCoreComponentPodConfig | indent 6 -}} -{{- end }} ---- -{{- end }} -apiVersion: apps/v1 -kind: Deployment -metadata: - name: mqtrigger-nats-streaming - labels: - chart: "{{ .Chart.Name }}-{{ .Chart.Version }}" - svc: mqtrigger - messagequeue: nats-streaming -spec: - replicas: 1 - selector: - matchLabels: - svc: mqtrigger - messagequeue: nats-streaming - template: - metadata: - labels: - svc: mqtrigger - messagequeue: nats-streaming - spec: - containers: - - name: mqtrigger - image: {{ include "fission-bundleImage" . | quote }} - imagePullPolicy: {{ .Values.pullPolicy }} - command: ["/fission-bundle"] - args: ["--mqt", "--routerUrl", "http://router.{{ .Release.Namespace }}"] - env: - - name: MESSAGE_QUEUE_TYPE - value: nats-streaming - - name: MESSAGE_QUEUE_CLUSTER_ID - value: {{ .Values.nats.clusterID }} - - name: MESSAGE_QUEUE_QUEUE_GROUP - value: {{ .Values.nats.queueGroup }} - - name: MESSAGE_QUEUE_CLIENT_ID - value: {{ .Values.nats.clientID }} - - name: MESSAGE_QUEUE_URL - {{- if .Values.nats.authToken }} - value: nats://{{ .Values.nats.authToken }}@{{ .Values.nats.hostaddress }} - {{- else }} - value: nats://{{ .Values.nats.hostaddress }} - {{- end }} - - name: DEBUG_ENV - value: {{ .Values.debugEnv | quote }} - - name: PPROF_ENABLED - value: {{ .Values.pprof.enabled | quote }} - {{- include "opentracing.envs" . | indent 8 }} - {{- include "opentelemtry.envs" . | indent 8 }} - {{- with .Values.imagePullSecrets }} - imagePullSecrets: - {{- toYaml . | nindent 8 }} - {{- end }} - serviceAccountName: fission-svc -{{- if .Values.extraCoreComponentPodConfig }} -{{ toYaml .Values.extraCoreComponentPodConfig | indent 6 -}} -{{- end }} -{{- end }} diff --git a/charts/fission-all/templates/mqt-fission-nats/svc.yaml b/charts/fission-all/templates/mqt-fission-nats/svc.yaml deleted file mode 100644 index 5129f43f..00000000 --- a/charts/fission-all/templates/mqt-fission-nats/svc.yaml +++ /dev/null @@ -1,19 +0,0 @@ -{{- if and .Values.nats.enabled (not .Values.nats.external) }} -apiVersion: v1 -kind: Service -metadata: - name: nats-streaming - labels: - svc: nats-streaming - chart: "{{ .Chart.Name }}-{{ .Chart.Version }}" -spec: - type: {{ .Values.serviceType }} - ports: - - port: 4222 - targetPort: 4222 -{{- if eq .Values.serviceType "NodePort" }} - nodePort: {{ .Values.natsStreamingPort }} -{{- end }} - selector: - svc: nats-streaming -{{- end }} \ No newline at end of file diff --git a/charts/fission-all/values.yaml b/charts/fission-all/values.yaml index c76abfd3..34407eab 100644 --- a/charts/fission-all/values.yaml +++ b/charts/fission-all/values.yaml @@ -320,51 +320,6 @@ timer: ## resources: {} -## Message queue trigger config -## NATS Streaming, enabled by default -## -nats: - ## whether or not to use NATS - ## - enabled: false - - ## if true, don't install NATS, but - ## use the existing NATS cluster - ## - external: false - - ## Address of NATS server (domain:port) - ## change from default for external NATS - ## - hostaddress: "nats-streaming:4222" - - ## Authorization token to use with NATS - ## - authToken: "defaultFissionAuthToken" - - ## NATS streaming clusterID - ## - clusterID: "fissionMQTrigger" - - ## Client name registered with NATS streaming - ## - clientID: "fission" - - ## Queue group registered with NATS streaming - ## - queueGroup: "fission-messageQueueNatsTrigger" - - ## The image to use for NATS streaming server - ## - streamingserver: - image: nats-streaming - tag: "0.23.0" - -## Port at which NATS streaming service should be exposed -## (only if nats enabled and not external) -## -natsStreamingPort: 31316 - ## Azure-storage-queue: enable and configure the details ## azureStorageQueue: diff --git a/cmd/fission-bundle/mqtrigger/mqtrigger.go b/cmd/fission-bundle/mqtrigger/mqtrigger.go index cdea92a9..14cb601d 100644 --- a/cmd/fission-bundle/mqtrigger/mqtrigger.go +++ b/cmd/fission-bundle/mqtrigger/mqtrigger.go @@ -33,7 +33,6 @@ import ( "github.com/fission/fission/pkg/mqtrigger/messageQueue" _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/azurequeuestorage" _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/kafka" - _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/nats" ) func Start(ctx context.Context, logger *zap.Logger, routerUrl string) error { diff --git a/cmd/fission-cli/app/app.go b/cmd/fission-cli/app/app.go index 7e7d721e..76a99d20 100644 --- a/cmd/fission-cli/app/app.go +++ b/cmd/fission-cli/app/app.go @@ -41,7 +41,6 @@ import ( "github.com/fission/fission/pkg/fission-cli/util" _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/azurequeuestorage" _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/kafka" - _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/nats" ) const ( diff --git a/go.mod b/go.mod index c3c77c34..adfc0f35 100644 --- a/go.mod +++ b/go.mod @@ -24,9 +24,6 @@ require ( github.com/influxdata/influxdb v1.9.6 github.com/mholt/archiver/v3 v3.5.1 github.com/minio/minio-go v6.0.14+incompatible - github.com/nats-io/nats-streaming-server v0.24.3 - github.com/nats-io/nats.go v1.13.1-0.20220308171302-2f2f6968e98d - github.com/nats-io/stan.go v0.10.2 github.com/ory/dockertest v3.3.5+incompatible github.com/pkg/errors v0.9.1 github.com/prometheus/client_golang v1.12.1 @@ -76,7 +73,6 @@ require ( github.com/PuerkitoBio/urlesc v0.0.0-20170810143723-de5bf2ad4578 // indirect github.com/acomagu/bufpipe v1.0.3 // indirect github.com/andybalholm/brotli v1.0.1 // indirect - github.com/armon/go-metrics v0.3.10 // indirect github.com/aws/aws-sdk-go v1.42.34 // indirect github.com/beorn7/perks v1.0.1 // indirect github.com/blend/go-sdk v1.20220212.5 // indirect @@ -117,9 +113,7 @@ require ( github.com/gotestyourself/gotestyourself v2.2.0+incompatible // indirect github.com/grpc-ecosystem/grpc-gateway v1.16.0 // indirect github.com/hashicorp/errwrap v1.0.0 // indirect - github.com/hashicorp/go-immutable-radix v1.3.1 // indirect github.com/hashicorp/go-uuid v1.0.2 // indirect - github.com/hashicorp/golang-lru v0.5.4 // indirect github.com/inconshreveable/mousetrap v1.0.0 // indirect github.com/jbenet/go-context v0.0.0-20150711004518-d14ea06fba99 // indirect github.com/jcmturner/aescts/v2 v2.0.0 // indirect @@ -133,6 +127,7 @@ require ( github.com/kevinburke/ssh_config v0.0.0-20201106050909-4977a11b4351 // indirect github.com/klauspost/compress v1.14.4 // indirect github.com/klauspost/pgzip v1.2.5 // indirect + github.com/lib/pq v1.10.4 // indirect github.com/mailru/easyjson v0.7.6 // indirect github.com/mattn/go-colorable v0.1.12 // indirect github.com/mattn/go-isatty v0.0.14 // indirect @@ -141,8 +136,6 @@ require ( github.com/moby/spdystream v0.2.0 // indirect github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect github.com/modern-go/reflect2 v1.0.2 // indirect - github.com/nats-io/nkeys v0.3.0 // indirect - github.com/nats-io/nuid v1.0.1 // indirect github.com/nwaples/rardecode v1.1.0 // indirect github.com/opencontainers/go-digest v1.0.0 // indirect github.com/opencontainers/image-spec v1.0.2 // indirect diff --git a/go.sum b/go.sum index f82044ce..b42371a5 100644 --- a/go.sum +++ b/go.sum @@ -81,8 +81,6 @@ github.com/Azure/go-autorest/tracing v0.6.0 h1:TYi4+3m5t6K48TGI9AUdb+IzbnSxvnvUM github.com/Azure/go-autorest/tracing v0.6.0/go.mod h1:+vhtPC754Xsa23ID7GlGsrdKBpUA79WCAKPPZVC2DeU= github.com/BurntSushi/toml v0.3.1/go.mod h1:xHWCNGjB5oqiDr8zfno3MHue2Ht5sIBksp03qcyfWMU= github.com/BurntSushi/xgb v0.0.0-20160522181843-27f122750802/go.mod h1:IVnqGOEym/WlBOVXweHU+Q+/VP0lqqI8lqeDx9IjBqo= -github.com/DataDog/datadog-go v2.2.0+incompatible/go.mod h1:LButxg5PwREeZtORoXG3tL4fMGNddJ+vMq1mwgfaqoQ= -github.com/DataDog/datadog-go v3.2.0+incompatible/go.mod h1:LButxg5PwREeZtORoXG3tL4fMGNddJ+vMq1mwgfaqoQ= github.com/Microsoft/go-winio v0.4.14/go.mod h1:qXqCSQ3Xa7+6tgxaGTIe4Kpcdsi+P8jBhyzoq1bpyYA= github.com/Microsoft/go-winio v0.4.16/go.mod h1:XB6nPKklQyQ7GC9LdcBEcBl8PF76WugXOPRXwdLnMv0= github.com/Microsoft/go-winio v0.5.1 h1:aPJp2QD7OOrhO5tQXqQoGSJc+DjDtWTGLOmNyAm6FgY= @@ -119,9 +117,6 @@ github.com/antlr/antlr4/runtime/Go/antlr v0.0.0-20210826220005-b48c857c3a0e/go.m github.com/armon/circbuf v0.0.0-20150827004946-bbbad097214e/go.mod h1:3U/XgcO3hCbHZ8TKRvWD2dDTCfh9M9ya+I9JpbB7O8o= github.com/armon/consul-api v0.0.0-20180202201655-eb2c6b5be1b6/go.mod h1:grANhF5doyWs3UAsr3K4I6qtAmlQcZDesFNEHPZAzj8= github.com/armon/go-metrics v0.0.0-20180917152333-f0300d1749da/go.mod h1:Q73ZrmVTwzkszR9V5SSuryQ31EELlFMUz1kKyl939pY= -github.com/armon/go-metrics v0.0.0-20190430140413-ec5e00d3c878/go.mod h1:3AMJUQhVx52RsWOnlkpikZr01T/yAVN2gn0861vByNg= -github.com/armon/go-metrics v0.3.10 h1:FR+drcQStOe+32sYyJYyZ7FIdgoGGBnwLl+flodp8Uo= -github.com/armon/go-metrics v0.3.10/go.mod h1:4O98XIr/9W0sxpJ8UaYkvjk10Iff7SnFrb4QAOwNTFc= github.com/armon/go-radix v0.0.0-20180808171621-7fddfc383310/go.mod h1:ufUuZ+zHj4x4TnLV4JWEpy2hxWSpsRywHrMgIH9cCH8= github.com/armon/go-socks5 v0.0.0-20160902184237-e75332964ef5 h1:0CwZNZbxp69SHPdPJAN/hZIm0C4OItdklCFmMRWYpio= github.com/armon/go-socks5 v0.0.0-20160902184237-e75332964ef5/go.mod h1:wHh0iHkYZB8zMSxRWpUBQtwG5a7fFgvEO+odwuTv2gs= @@ -160,8 +155,6 @@ github.com/chzyer/logex v1.1.10/go.mod h1:+Ywpsq7O8HXn0nuIou7OrIPyXbp3wmkHB+jjWR github.com/chzyer/readline v0.0.0-20180603132655-2972be24d48e/go.mod h1:nSuG5e5PlCu98SY8svDHJxuZscDgtXS6KTTbou5AhLI= github.com/chzyer/test v0.0.0-20180213035817-a1ea475d72b1/go.mod h1:Q3SI9o4m/ZMnBNeIyt5eFwwo7qiLfzFZmjNmxjkiQlU= github.com/cilium/ebpf v0.7.0/go.mod h1:/oI2+1shJiTGAMgl6/RgJr36Eo1jzrRcAWbcXO2usCA= -github.com/circonus-labs/circonus-gometrics v2.3.1+incompatible/go.mod h1:nmEj6Dob7S7YxXgwXpfOuvO54S+tGdZdw9fuRZt25Ag= -github.com/circonus-labs/circonusllhist v0.1.3/go.mod h1:kMXHVDlOchFAehlya5ePtbp5jckzBHf4XRpQvBOLI+I= github.com/client9/misspell v0.3.4/go.mod h1:qj6jICC3Q7zFZvVWo7KLAzC3yx5G7kyvSDkc90ppPyw= github.com/cncf/udpa/go v0.0.0-20191209042840-269d4d468f6f/go.mod h1:M8M6+tZqaGXZJjfX53e64911xZQV5JYwmTeXPW+k8Sc= github.com/cncf/udpa/go v0.0.0-20200629203442-efcf912fb354/go.mod h1:WmhPx2Nbnhtbo57+VJT5O0JRkEi1Wbu0z5j0R8u5Hbk= @@ -316,7 +309,6 @@ github.com/go-openapi/swag v0.19.5/go.mod h1:POnQmlKehdgb5mhVOsnJFsivZCEZ/vjK9gh github.com/go-openapi/swag v0.19.14/go.mod h1:QYRuS/SOXUCsnplDa677K7+DxSOj6IPNl/eQntq43wQ= github.com/go-openapi/swag v0.19.15 h1:D2NRCBzS9/pEY3gP9Nl8aDqGUcPFrwG2p+CNFrLyrCM= github.com/go-openapi/swag v0.19.15/go.mod h1:QYRuS/SOXUCsnplDa677K7+DxSOj6IPNl/eQntq43wQ= -github.com/go-sql-driver/mysql v1.6.0/go.mod h1:DCzpHaOWr8IXmIStZouvnhqoel9Qv2LBy8hT2VhHyBg= github.com/go-stack/stack v1.8.0/go.mod h1:v0f6uXyyMGvRgIKkXu+yp6POWl0qKG85gN/melR3HDY= github.com/go-task/slim-sprig v0.0.0-20210107165309-348f09dbbbc0/go.mod h1:fyg7847qk6SyHyPtNmDHnmrv/HOrqktSC+C9fM+CJOE= github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= @@ -445,22 +437,12 @@ github.com/hashicorp/consul/api v1.1.0/go.mod h1:VmuI/Lkw1nC05EYQWNKwWGbkg+FbDBt github.com/hashicorp/consul/sdk v0.1.1/go.mod h1:VKf9jXwCTEY1QZP2MOLRhb5i/I/ssyNV1vwHyQBF0x8= github.com/hashicorp/errwrap v1.0.0 h1:hLrqtEDnRye3+sgx6z4qVLNuviH3MR5aQ0ykNJa/UYA= github.com/hashicorp/errwrap v1.0.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4= -github.com/hashicorp/go-cleanhttp v0.5.0/go.mod h1:JpRdi6/HCYpAwUzNwuwqhbovhLtngrth3wmdIIUrZ80= github.com/hashicorp/go-cleanhttp v0.5.1/go.mod h1:JpRdi6/HCYpAwUzNwuwqhbovhLtngrth3wmdIIUrZ80= -github.com/hashicorp/go-hclog v0.9.1/go.mod h1:5CU+agLiy3J7N7QjHK5d05KxGsuXiQLrjA0H7acj2lQ= -github.com/hashicorp/go-hclog v1.1.0 h1:QsGcniKx5/LuX2eYoeL+Np3UKYPNaN7YKpTh29h8rbw= -github.com/hashicorp/go-hclog v1.1.0/go.mod h1:whpDNt7SSdeAju8AWKIWsul05p54N/39EeqMAyrmvFQ= github.com/hashicorp/go-immutable-radix v1.0.0/go.mod h1:0y9vanUI8NX6FsYoO3zeMjhV/C5i9g4Q3DwcSNZ4P60= -github.com/hashicorp/go-immutable-radix v1.3.1 h1:DKHmCUm2hRBK510BaiZlwvpD40f8bJFeZnpfm2KLowc= -github.com/hashicorp/go-immutable-radix v1.3.1/go.mod h1:0y9vanUI8NX6FsYoO3zeMjhV/C5i9g4Q3DwcSNZ4P60= github.com/hashicorp/go-msgpack v0.5.3/go.mod h1:ahLV/dePpqEmjfWmKiqvPkv/twdG7iPBM1vqhUKIvfM= -github.com/hashicorp/go-msgpack v0.5.5/go.mod h1:ahLV/dePpqEmjfWmKiqvPkv/twdG7iPBM1vqhUKIvfM= -github.com/hashicorp/go-msgpack v1.1.5 h1:9byZdVjKTe5mce63pRVNP1L7UAmdHOTEMGehn6KvJWs= -github.com/hashicorp/go-msgpack v1.1.5/go.mod h1:gWVc3sv/wbDmR3rQsj1CAktEZzoz1YNK9NfGLXJ69/4= github.com/hashicorp/go-multierror v1.0.0/go.mod h1:dHtQlpGsu+cZNNAkkCN/P3hoUDHhCYQXV3UM06sGGrk= github.com/hashicorp/go-multierror v1.1.1 h1:H5DkEtf6CXdFp0N0Em5UCwQpXMWke8IA0+lD48awMYo= github.com/hashicorp/go-multierror v1.1.1/go.mod h1:iw975J/qwKPdAO1clOe2L8331t/9/fmwbPZ6JB6eMoM= -github.com/hashicorp/go-retryablehttp v0.5.3/go.mod h1:9B5zBasrRhHXnJnui7y6sL7es7NDiJgTc6Er0maI1Xs= github.com/hashicorp/go-rootcerts v1.0.0/go.mod h1:K6zTfqpRlCUIjkwsN4Z+hiSfzSTQa6eBIzfwKfwNnHU= github.com/hashicorp/go-sockaddr v1.0.0/go.mod h1:7Xibr9yA9JjQq1JpNB2Vw7kxv8xerXegt+ozgdvDeDU= github.com/hashicorp/go-syslog v1.0.0/go.mod h1:qPfqrKkXGihmCqbJM2mZgkZGvKG1dFdvsLplgctolz4= @@ -471,14 +453,10 @@ github.com/hashicorp/go-uuid v1.0.2/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/b github.com/hashicorp/go.net v0.0.1/go.mod h1:hjKkEWcCURg++eb33jQU7oqQcI9XDCnUzHA0oac0k90= github.com/hashicorp/golang-lru v0.5.0/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8= github.com/hashicorp/golang-lru v0.5.1/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8= -github.com/hashicorp/golang-lru v0.5.4 h1:YDjusn29QI/Das2iO9M0BHnIbxPeyuCHsjMW+lJfyTc= -github.com/hashicorp/golang-lru v0.5.4/go.mod h1:iADmTwqILo4mZ8BN3D2Q6+9jd8WM5uGBxy+E8yxSoD4= github.com/hashicorp/hcl v1.0.0/go.mod h1:E5yfLk+7swimpb2L/Alb/PJmXilQ/rhwaUYs4T20WEQ= github.com/hashicorp/logutils v1.0.0/go.mod h1:QIAnNjmIWmVIIkWDTG1z5v++HQmx9WQRO+LraFDTW64= github.com/hashicorp/mdns v1.0.0/go.mod h1:tL+uN++7HEJ6SQLQ2/p+z2pH24WQKWjBPkE0mNTz8vQ= github.com/hashicorp/memberlist v0.1.3/go.mod h1:ajVTdAv/9Im8oMAAj5G31PhhMCZJV2pPBoIllUwCN7I= -github.com/hashicorp/raft v1.3.6 h1:v5xW5KzByoerQlN/o31VJrFNiozgzGyDoMgDJgXpsto= -github.com/hashicorp/raft v1.3.6/go.mod h1:4Ak7FSPnuvmb0GV6vgIAJ4vYT4bek9bb6Q+7HVbyzqM= github.com/hashicorp/serf v0.8.2/go.mod h1:6hOLApaqBFA1NXqRQAsxw9QxuDEvNxSQRwA/JwenrHc= github.com/hpcloud/tail v1.0.0/go.mod h1:ab1qPbhIpdTxEkNHXyeSf5vhxWSCs/tWer42PpOxQnU= github.com/ianlancetaylor/demangle v0.0.0-20181102032728-5e5cf60278f6/go.mod h1:aSSvb/t6k1mPoxDqO4vJh6VOCGPwU4O0C2/Eqndh1Sc= @@ -518,7 +496,6 @@ github.com/josharian/intern v1.0.0/go.mod h1:5DoeVV0s6jJacbCEi61lwdGj/aVlrQvzHFF github.com/jpillora/backoff v1.0.0 h1:uvFg412JmmHBHw7iwprIxkPMI+sGQ4kzOWsMeHnm2EA= github.com/jpillora/backoff v1.0.0/go.mod h1:J/6gKK9jxlEcS3zixgDgUAsiuZ7yrSoa/FX5e0EB2j4= github.com/json-iterator/go v1.1.6/go.mod h1:+SdeFBvtyEkXs7REEP0seUULqWtbJapLOCVDaaPEHmU= -github.com/json-iterator/go v1.1.9/go.mod h1:KdQUCv79m/52Kvf8AW2vK1V8akMuk1QjK/uOdHXbAo4= github.com/json-iterator/go v1.1.10/go.mod h1:KdQUCv79m/52Kvf8AW2vK1V8akMuk1QjK/uOdHXbAo4= github.com/json-iterator/go v1.1.11/go.mod h1:KdQUCv79m/52Kvf8AW2vK1V8akMuk1QjK/uOdHXbAo4= github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= @@ -567,13 +544,10 @@ github.com/mailru/easyjson v0.7.6/go.mod h1:xzfreul335JAWq5oZzymOObrkdz5UnU4kGfJ github.com/matryer/is v1.2.0 h1:92UTHpy8CDwaJ08GqLDzhhuixiBUUD1p3AU6PHddz4A= github.com/matryer/is v1.2.0/go.mod h1:2fLPjFQM9rhQ15aVEtbuwhJinnOqrmgXPNdZsdwlWXA= github.com/mattn/go-colorable v0.0.9/go.mod h1:9vuHe8Xs5qXnSaW/c/ABM9alt+Vo+STaOChaDxuIBZU= -github.com/mattn/go-colorable v0.1.4/go.mod h1:U0ppj6V5qS13XJ6of8GYAs25YV2eR4EVcfRqFIhoBtE= github.com/mattn/go-colorable v0.1.9/go.mod h1:u6P/XSegPjTcexA+o6vUJrdnUu04hMope9wVRipJSqc= github.com/mattn/go-colorable v0.1.12 h1:jF+Du6AlPIjs2BiUiQlKOX0rt3SujHxPnksPKZbaA40= github.com/mattn/go-colorable v0.1.12/go.mod h1:u5H1YNBxpqRaxsYJYSkiCWKzEfiAb1Gb520KVy5xxl4= github.com/mattn/go-isatty v0.0.3/go.mod h1:M+lRXTBqGeGNdLjl/ufCoiOlB5xdOkqRJdNxMWT7Zi4= -github.com/mattn/go-isatty v0.0.8/go.mod h1:Iq45c/XA43vh69/j3iqttzPXn0bhXyGjM0Hdxcsrc5s= -github.com/mattn/go-isatty v0.0.10/go.mod h1:qgIWMr58cqv1PHHyhnkY9lrL7etaEgOFcMEpPG5Rm84= github.com/mattn/go-isatty v0.0.12/go.mod h1:cbi8OIDigv2wuxKPP5vlRcQ1OAZbq2CE4Kysco4FUpU= github.com/mattn/go-isatty v0.0.14 h1:yVuAays6BHfxijgZPzw+3Zlu5yQgKGP2/hcQbHb7S9Y= github.com/mattn/go-isatty v0.0.14/go.mod h1:7GGIvUiUoEMVVmxf/4nioHXj79iQHKdU27kJ6hsGG94= @@ -583,8 +557,6 @@ github.com/matttproud/golang_protobuf_extensions v1.0.2-0.20181231171920-c182aff github.com/mholt/archiver/v3 v3.5.1 h1:rDjOBX9JSF5BvoJGvjqK479aL70qh9DIpZCl+k7Clwo= github.com/mholt/archiver/v3 v3.5.1/go.mod h1:e3dqJ7H78uzsRSEACH1joayhuSyhnonssnDhppzS1L4= github.com/miekg/dns v1.0.14/go.mod h1:W1PPwlIAgtquWBMBEV9nkV9Cazfe8ScdGz/Lj7v3Nrg= -github.com/minio/highwayhash v1.0.2 h1:Aak5U0nElisjDCfPSG79Tgzkn2gl66NxOMspRrKnA/g= -github.com/minio/highwayhash v1.0.2/go.mod h1:BQskDq+xkJ12lmlUUi7U0M5Swg3EWR+dLTk+kldvVxY= github.com/minio/minio-go v6.0.14+incompatible h1:fnV+GD28LeqdN6vT2XdGKW8Qe/IfjJDswNVuni6km9o= github.com/minio/minio-go v6.0.14+incompatible/go.mod h1:7guKYtitv8dktvNUGrhzmNlA5wrAABTQXCoesZdFQO8= github.com/mitchellh/cli v1.0.0/go.mod h1:hNIlj7HEI86fIcpObd7a0FcrxTWetlwJDGcceTlRvqc= @@ -616,21 +588,6 @@ github.com/mwitkow/go-conntrack v0.0.0-20161129095857-cc309e4a2223/go.mod h1:qRW github.com/mwitkow/go-conntrack v0.0.0-20190716064945-2f068394615f h1:KUppIJq7/+SVif2QVs3tOP0zanoHgBEVAwHxUSIzRqU= github.com/mwitkow/go-conntrack v0.0.0-20190716064945-2f068394615f/go.mod h1:qRWi+5nqEBWmkhHvq77mSJWrCKwh8bxhgT7d/eI7P4U= github.com/mxk/go-flowrate v0.0.0-20140419014527-cca7078d478f/go.mod h1:ZdcZmHo+o7JKHSa8/e818NopupXU1YMK5fe1lsApnBw= -github.com/nats-io/jwt/v2 v2.2.1-0.20220113022732-58e87895b296 h1:vU9tpM3apjYlLLeY23zRWJ9Zktr5jp+mloR942LEOpY= -github.com/nats-io/jwt/v2 v2.2.1-0.20220113022732-58e87895b296/go.mod h1:0tqz9Hlu6bCBFLWAASKhE5vUA4c24L9KPUUgvwumE/k= -github.com/nats-io/nats-server/v2 v2.7.4 h1:c+BZJ3rGzUKCBIM4IXO8uNT2u1vajGbD1kPA6wqCEaM= -github.com/nats-io/nats-server/v2 v2.7.4/go.mod h1:1vZ2Nijh8tcyNe8BDVyTviCd9NYzRbubQYiEHsvOQWc= -github.com/nats-io/nats-streaming-server v0.24.3 h1:uZez8jBkXscua++jaDsK7DhpSAkizdetar6yWbPMRco= -github.com/nats-io/nats-streaming-server v0.24.3/go.mod h1:rqWfyCbxlhKj//fAp8POdQzeADwqkVhZcoWlbhkuU5w= -github.com/nats-io/nats.go v1.13.0/go.mod h1:BPko4oXsySz4aSWeFgOHLZs3G4Jq4ZAyE6/zMCxRT6w= -github.com/nats-io/nats.go v1.13.1-0.20220308171302-2f2f6968e98d h1:zJf4l8Kp67RIZhoVeniSLZs69SHNgjLHz0aNsqPPlx8= -github.com/nats-io/nats.go v1.13.1-0.20220308171302-2f2f6968e98d/go.mod h1:BPko4oXsySz4aSWeFgOHLZs3G4Jq4ZAyE6/zMCxRT6w= -github.com/nats-io/nkeys v0.3.0 h1:cgM5tL53EvYRU+2YLXIK0G2mJtK12Ft9oeooSZMA2G8= -github.com/nats-io/nkeys v0.3.0/go.mod h1:gvUNGjVcM2IPr5rCsRsC6Wb3Hr2CQAm08dsxtV6A5y4= -github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw= -github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= -github.com/nats-io/stan.go v0.10.2 h1:gQLd05LhzmhFkHm3/qP/klYHfM/hys45GyHa1Uly/kI= -github.com/nats-io/stan.go v0.10.2/go.mod h1:vo2ax8K2IxaR3JtEMLZRFKIdoK/3o1/PKueapB7ezX0= github.com/ncw/swift v1.0.49/go.mod h1:23YIA4yWVnGwv2dQlN4bB7egfYX6YLn0Yo/S6zZO/ZM= github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLANZcx1PVRCS0qkT7pwLkGfwJo4zjcN/Tysno= github.com/nwaples/rardecode v1.1.0 h1:vSxaY8vQhOcVr4mm5e8XllHWTiM4JF507A0Katqw7MQ= @@ -662,7 +619,6 @@ github.com/opentracing/opentracing-go v1.1.0/go.mod h1:UkNAQd3GIcIGf0SeVgPpRdFSt github.com/ory/dockertest v3.3.5+incompatible h1:iLLK6SQwIhcbrG783Dghaaa3WPzGc+4Emza6EbVUUGA= github.com/ory/dockertest v3.3.5+incompatible/go.mod h1:1vX4m9wsvi00u5bseYwXaSnhNrne+V0E6LAcBILJdPs= github.com/pascaldekloe/goe v0.0.0-20180627143212-57f6aae5913c/go.mod h1:lzWF7FIEvWOWxwDKqyGYQf6ZUaNfKdP144TG7ZOy1lc= -github.com/pascaldekloe/goe v0.1.0/go.mod h1:lzWF7FIEvWOWxwDKqyGYQf6ZUaNfKdP144TG7ZOy1lc= github.com/pelletier/go-toml v1.2.0/go.mod h1:5z9KED0ma1S8pY6P1sdut58dfprrGBbd/94hg7ilaic= github.com/pelletier/go-toml v1.9.3/go.mod h1:u1nR/EPcESfeI/szUZKdtJ0xRNbUoANCkoOuaOx1Y+c= github.com/peterbourgon/diskv v2.0.1+incompatible/go.mod h1:uqqh8zWWbv1HBMNONnaR/tNboyR3/BZd58JJSHlUSCU= @@ -683,10 +639,8 @@ github.com/posener/complete v1.1.1/go.mod h1:em0nMJCgc9GFtwrmVmEMR/ZL6WyhyjMBndr github.com/pquerna/cachecontrol v0.0.0-20171018203845-0dec1b30a021/go.mod h1:prYjPmNq4d1NPVmpShWobRqXY3q7Vp+80DqgxxUrUIA= github.com/pquerna/ffjson v0.0.0-20190813045741-dac163c6c0a9/go.mod h1:YARuvh7BUWHNhzDq2OM5tzR2RiCcN2D7sapiKyCel/M= github.com/prometheus/client_golang v0.9.1/go.mod h1:7SWBe2y4D6OKWSNQJUaRYU/AaXPKyh/dDVn+NZz0KFw= -github.com/prometheus/client_golang v0.9.2/go.mod h1:OsXs2jCmiKlQ1lTBmv21f2mNfw4xf/QclQDMrYNZzcM= github.com/prometheus/client_golang v0.9.3/go.mod h1:/TN21ttK/J9q6uSwhBd54HahCDft0ttaMvbicHlPoso= github.com/prometheus/client_golang v1.0.0/go.mod h1:db9x61etRT2tGnBNRi70OPL5FsnadC4Ky3P0J6CfImo= -github.com/prometheus/client_golang v1.4.0/go.mod h1:e9GMxYsXl05ICDXkRhurwBS4Q3OK1iX/F2sw+iXX5zU= github.com/prometheus/client_golang v1.7.1/go.mod h1:PY5Wy2awLA44sXw4AOSfFBetzPP4j5+D6mVACh+pe2M= github.com/prometheus/client_golang v1.11.0/go.mod h1:Z6t4BnS23TR94PD6BsDNk8yVqroYurpAkEiz0P2BEV0= github.com/prometheus/client_golang v1.12.1 h1:ZiaPsmm9uiBeaSMRznKsCDNtPCS0T3JVDGF+06gjBzk= @@ -697,20 +651,16 @@ github.com/prometheus/client_model v0.0.0-20190812154241-14fe0d1b01d4/go.mod h1: github.com/prometheus/client_model v0.2.0 h1:uq5h0d+GuxiXLJLNABMgp2qUWDPiLvgCzz2dUR+/W/M= github.com/prometheus/client_model v0.2.0/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA= github.com/prometheus/common v0.0.0-20181113130724-41aa239b4cce/go.mod h1:daVV7qP5qjZbuso7PdcryaAu0sAZbrN9i7WWcTMWvro= -github.com/prometheus/common v0.0.0-20181126121408-4724e9255275/go.mod h1:daVV7qP5qjZbuso7PdcryaAu0sAZbrN9i7WWcTMWvro= github.com/prometheus/common v0.4.0/go.mod h1:TNfzLD0ON7rHzMJeJkieUDPYmFC7Snx/y86RQel1bk4= github.com/prometheus/common v0.4.1/go.mod h1:TNfzLD0ON7rHzMJeJkieUDPYmFC7Snx/y86RQel1bk4= -github.com/prometheus/common v0.9.1/go.mod h1:yhUN8i9wzaXS3w1O07YhxHEBxD+W35wd8bs7vj7HSQ4= github.com/prometheus/common v0.10.0/go.mod h1:Tlit/dnDKsSWFlCLTWaA1cyBgKHSMdTB80sz/V91rCo= github.com/prometheus/common v0.26.0/go.mod h1:M7rCNAaPfAosfx8veZJCuw84e35h3Cfd9VFqTh1DIvc= github.com/prometheus/common v0.28.0/go.mod h1:vu+V0TpY+O6vW9J44gczi3Ap/oXXR10b+M/gUGO4Hls= github.com/prometheus/common v0.32.1 h1:hWIdL3N2HoUx3B8j3YN9mWor0qhY/NlEKZEaXxuIRh4= github.com/prometheus/common v0.32.1/go.mod h1:vu+V0TpY+O6vW9J44gczi3Ap/oXXR10b+M/gUGO4Hls= github.com/prometheus/procfs v0.0.0-20181005140218-185b4288413d/go.mod h1:c3At6R/oaqEKCNdg8wHV1ftS6bRYblBhIjjI8uT2IGk= -github.com/prometheus/procfs v0.0.0-20181204211112-1dc9a6cbc91a/go.mod h1:c3At6R/oaqEKCNdg8wHV1ftS6bRYblBhIjjI8uT2IGk= github.com/prometheus/procfs v0.0.0-20190507164030-5867b95ac084/go.mod h1:TjEm7ze935MbeOT/UhFTIMYKhuLP4wbCsTZCD3I8kEA= github.com/prometheus/procfs v0.0.2/go.mod h1:TjEm7ze935MbeOT/UhFTIMYKhuLP4wbCsTZCD3I8kEA= -github.com/prometheus/procfs v0.0.8/go.mod h1:7Qr8sr6344vo1JqZ6HhLceV9o3AJ1Ff+GxbHq6oeK9A= github.com/prometheus/procfs v0.1.3/go.mod h1:lV6e/gmhEcM9IjHGsFOCxxuZ+z1YqCvr4OA4YeYWdaU= github.com/prometheus/procfs v0.6.0/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1xBZuNvfVA= github.com/prometheus/procfs v0.7.3 h1:4jVXhlkAyzOScmCkXBTOLRLTz8EeU+eyjrwB/EPq0VU= @@ -786,7 +736,6 @@ github.com/subosito/gotenv v1.2.0/go.mod h1:N0PQaV/YGNqwC0u51sEeR/aUtSLEXKX9iv69 github.com/syndtr/gocapability v0.0.0-20200815063812-42c35b437635/go.mod h1:hkRG7XYTFWNJGYcbNJQlaLq0fg1yr4J4t/NcTQtrfww= github.com/tmc/grpc-websocket-proxy v0.0.0-20190109142713-0ad062ec5ee5/go.mod h1:ncp9v5uamzpCO7NfCPTXjqaC+bZgJeR0sMTm6dMHP7U= github.com/tmc/grpc-websocket-proxy v0.0.0-20201229170055-e5319fda7802/go.mod h1:ncp9v5uamzpCO7NfCPTXjqaC+bZgJeR0sMTm6dMHP7U= -github.com/tv42/httpunix v0.0.0-20150427012821-b75d8614f926/go.mod h1:9ESjWnEqriFuLhtthL60Sar/7RFoluCcXsuvEwTV5KM= github.com/tv42/httpunix v0.0.0-20191220191345-2ba4b9c3382c/go.mod h1:hzIxponao9Kjc7aWznkXaL4U4TWaDSs8zcsY4Ka08nM= github.com/uber/jaeger-client-go v2.25.0+incompatible/go.mod h1:WVhlPFC8FDjOFMMWRy2pZqQJSXxYSwNYOkTr/Z6d3Kk= github.com/uber/jaeger-client-go v2.28.0+incompatible h1:G4QSBfvPKvg5ZM2j9MrJFdfI5iSljY/WnJqOGFao6HI= @@ -817,7 +766,6 @@ github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9dec github.com/yuin/goldmark v1.3.5/go.mod h1:mwnBkeHKe2W/ZEtQ+71ViKU8L12m81fl3OWwC1Zlc8k= github.com/yuin/goldmark v1.4.0/go.mod h1:mwnBkeHKe2W/ZEtQ+71ViKU8L12m81fl3OWwC1Zlc8k= go.etcd.io/bbolt v1.3.2/go.mod h1:IbVyRI1SCnLcuJnV2u8VeU0CEYM7e686BmAb1XKL+uU= -go.etcd.io/bbolt v1.3.6 h1:/ecaJf0sk1l4l6V4awd65v2C3ILy7MSj+s/x1ADCIMU= go.etcd.io/bbolt v1.3.6/go.mod h1:qXsaaIqmgQH0T+OPdb99Bf+PKfBBQVAdyD6TY9G8XM4= go.etcd.io/etcd/api/v3 v3.5.0/go.mod h1:cbVKeC6lCfl7j/8jBhAK6aIYO9XOjdptoxU/nLQcPvs= go.etcd.io/etcd/client/pkg/v3 v3.5.0/go.mod h1:IJHfcCEKxYu1Os13ZdwCwIUTUVGYTSAM3YSwc9/Ac1g= @@ -905,11 +853,9 @@ golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPh golang.org/x/crypto v0.0.0-20201002170205-7f63de1d35b0/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= golang.org/x/crypto v0.0.0-20201112155050-0c6587e931a9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= golang.org/x/crypto v0.0.0-20210220033148-5ea612d1eb83/go.mod h1:jdWPYTVW3xRLrWPugEBEK3UY2ZEsg3UU495nc5E+M+I= -golang.org/x/crypto v0.0.0-20210314154223-e6e6c4f2bb5b/go.mod h1:T9bdIzuCu7OtxOm1hfPfRQxPLYneinmdGuTeoZ9dtd4= golang.org/x/crypto v0.0.0-20210322153248-0c34fe9e7dc2/go.mod h1:T9bdIzuCu7OtxOm1hfPfRQxPLYneinmdGuTeoZ9dtd4= golang.org/x/crypto v0.0.0-20210421170649-83a5a9bb288b/go.mod h1:T9bdIzuCu7OtxOm1hfPfRQxPLYneinmdGuTeoZ9dtd4= golang.org/x/crypto v0.0.0-20210817164053-32db794688a5/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc= -golang.org/x/crypto v0.0.0-20220112180741-5e0467b6c7ce/go.mod h1:IxCIyHEi3zRg3s0A5j5BB6A9Jmi73HwBIUl50j+osU4= golang.org/x/crypto v0.0.0-20220214200702-86341886e292/go.mod h1:IxCIyHEi3zRg3s0A5j5BB6A9Jmi73HwBIUl50j+osU4= golang.org/x/crypto v0.0.0-20220307211146-efcb8507fb70 h1:syTAU9FwmvzEoIYMqcPHOcVm4H3U5u90WsvuYgwpETU= golang.org/x/crypto v0.0.0-20220307211146-efcb8507fb70/go.mod h1:IxCIyHEi3zRg3s0A5j5BB6A9Jmi73HwBIUl50j+osU4= @@ -1047,9 +993,7 @@ golang.org/x/sys v0.0.0-20180909124046-d0be0721c37e/go.mod h1:STP8DvDyc/dI5b8T5h golang.org/x/sys v0.0.0-20181026203630-95b1ffbd15a5/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20181107165924-66b7b1311ac8/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20181116152217-5ac8a444bdc5/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= -golang.org/x/sys v0.0.0-20190130150945-aca44879d564/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= -golang.org/x/sys v0.0.0-20190222072716-a9d3bda3a223/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20190312061237-fead79001313/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20190422165155-953cdadca894/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= @@ -1063,7 +1007,6 @@ golang.org/x/sys v0.0.0-20190904154756-749cb33beabd/go.mod h1:h1NjWce9XRLGQEsW7w golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20191001151750-bb3f8db39f24/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20191005200804-aed5e4c7ecf9/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20191008105621-543471e840be/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20191026070338-33540a1f6037/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20191115151921-52ab43148777/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20191120155948-bd437916bb0e/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= @@ -1129,7 +1072,6 @@ golang.org/x/sys v0.0.0-20211116061358-0a5406a5449c/go.mod h1:oPkhp1MJrh7nUepCBc golang.org/x/sys v0.0.0-20211124211545-fe61309f8881/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20211205182925-97ca703d548d/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20211216021012-1d35b9e2eb4e/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.0.0-20220111092808-5a964db01320/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220114195835-da31bd327af9/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220307203707-22a9840ba4d7 h1:8IVLkfbr2cLhv0a/vKq4UFUcJym8RmDoDboxCFWEjYE= golang.org/x/sys v0.0.0-20220307203707-22a9840ba4d7/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= @@ -1165,7 +1107,6 @@ golang.org/x/tools v0.0.0-20190311212946-11955173bddd/go.mod h1:LCzVGOaR6xXOjkQ3 golang.org/x/tools v0.0.0-20190312151545-0bb0c0a6e846/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs= golang.org/x/tools v0.0.0-20190312170243-e65039ee4138/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs= golang.org/x/tools v0.0.0-20190328211700-ab21143f2384/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs= -golang.org/x/tools v0.0.0-20190424220101-1e8e1cfdf96b/go.mod h1:RgjU9mgBXZiqYHBnxXauZ1Gv1EHHAz9KjViQ78xBX0Q= golang.org/x/tools v0.0.0-20190425150028-36563e24a262/go.mod h1:RgjU9mgBXZiqYHBnxXauZ1Gv1EHHAz9KjViQ78xBX0Q= golang.org/x/tools v0.0.0-20190506145303-2d16b83fe98c/go.mod h1:RgjU9mgBXZiqYHBnxXauZ1Gv1EHHAz9KjViQ78xBX0Q= golang.org/x/tools v0.0.0-20190524140312-2c0ae7006135/go.mod h1:RgjU9mgBXZiqYHBnxXauZ1Gv1EHHAz9KjViQ78xBX0Q= diff --git a/pkg/apis/core/v1/const.go b/pkg/apis/core/v1/const.go index ae696d0d..922c164e 100644 --- a/pkg/apis/core/v1/const.go +++ b/pkg/apis/core/v1/const.go @@ -76,7 +76,6 @@ const ( ) const ( - MessageQueueTypeNats = "nats-streaming" MessageQueueTypeASQ = "azure-storage-queue" MessageQueueTypeKafka = "kafka" ) diff --git a/pkg/fission-cli/cmd/support/dump.go b/pkg/fission-cli/cmd/support/dump.go index 58147355..9b10a78d 100644 --- a/pkg/fission-cli/cmd/support/dump.go +++ b/pkg/fission-cli/cmd/support/dump.go @@ -83,15 +83,15 @@ func (opts *DumpSubCommand) do(input cli.Input) error { // fission component logs & spec "fission-components-svc-spec": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesService, - "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, nats-streaming, router, storagesvc, timer)"), + "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, router, storagesvc, timer)"), "fission-components-deployment-spec": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesDeployment, - "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, nats-streaming, router, storagesvc, timer)"), + "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, router, storagesvc, timer)"), "fission-components-daemonset-spec": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesDaemonSet, - "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, nats-streaming, router, storagesvc, timer)"), + "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, router, storagesvc, timer)"), "fission-components-pod-spec": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesPod, - "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, nats-streaming, router, storagesvc, timer)"), + "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, router, storagesvc, timer)"), "fission-components-pod-log": resources.NewKubernetesPodLogDumper(k8sClient, - "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, nats-streaming, router, storagesvc, timer)"), + "svc in (buildermgr, controller, executor, influxdb, kubewatcher, logger, mqtrigger, router, storagesvc, timer)"), // fission builder logs & spec "fission-builder-svc-spec": resources.NewKubernetesObjectDumper(k8sClient, resources.KubernetesService, "owner=buildermgr"), diff --git a/pkg/fission-cli/cmd/support/resources/crd.go b/pkg/fission-cli/cmd/support/resources/crd.go index 97028d66..f8d874fd 100644 --- a/pkg/fission-cli/cmd/support/resources/crd.go +++ b/pkg/fission-cli/cmd/support/resources/crd.go @@ -113,7 +113,7 @@ func (res CrdDumper) Dump(dumpDir string) { case CrdMessageQueueTrigger: var triggers []fv1.MessageQueueTrigger - for _, mqType := range []string{fv1.MessageQueueTypeNats, fv1.MessageQueueTypeASQ, fv1.MessageQueueTypeKafka} { + for _, mqType := range []string{fv1.MessageQueueTypeASQ, fv1.MessageQueueTypeKafka} { l, err := res.client.V1().MessageQueueTrigger().List(mqType, metav1.NamespaceAll) if err != nil { console.Warn(fmt.Sprintf("Error getting %v list: %v", res.crdType, err)) diff --git a/pkg/fission-cli/flag/flag.go b/pkg/fission-cli/flag/flag.go index 05a813bc..cad3ac46 100644 --- a/pkg/fission-cli/flag/flag.go +++ b/pkg/fission-cli/flag/flag.go @@ -155,7 +155,7 @@ var ( MqtName = Flag{Type: String, Name: flagkey.MqtName, Usage: "Message queue trigger name"} MqtFnName = Flag{Type: String, Name: flagkey.MqtFnName, Usage: "Function name"} - MqtMQType = Flag{Type: String, Name: flagkey.MqtMQType, Usage: "For mqtype \"fission\" => nats-streaming, azure-storage-queue, kafka\n\t\t\t\t\t For mqtype \"keda\" => kafka, aws-sqs-queue, aws-kinesis-stream, gcp-pubsub, stan, rabbitmq, redis", DefaultValue: "kafka"} + MqtMQType = Flag{Type: String, Name: flagkey.MqtMQType, Usage: "For mqtype \"fission\" => azure-storage-queue, kafka\n\t\t\t\t\t For mqtype \"keda\" => kafka, aws-sqs-queue, aws-kinesis-stream, gcp-pubsub, stan, rabbitmq, redis", DefaultValue: "kafka"} MqtTopic = Flag{Type: String, Name: flagkey.MqtTopic, Usage: "Message queue Topic the trigger listens on"} MqtRespTopic = Flag{Type: String, Name: flagkey.MqtRespTopic, Usage: "Topic that the function response is sent on (response discarded if unspecified)"} MqtErrorTopic = Flag{Type: String, Name: flagkey.MqtErrorTopic, Usage: "Topic that the function error messages are sent to (errors discarded if unspecified"} diff --git a/pkg/mqtrigger/messageQueue/nats/nats.go b/pkg/mqtrigger/messageQueue/nats/nats.go deleted file mode 100644 index 3158b8c2..00000000 --- a/pkg/mqtrigger/messageQueue/nats/nats.go +++ /dev/null @@ -1,245 +0,0 @@ -/* -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 nats - -import ( - "bytes" - "fmt" - "io" - "net/http" - "os" - "strings" - - nsUtil "github.com/nats-io/nats-streaming-server/util" - ns "github.com/nats-io/stan.go" - "go.uber.org/zap" - - fv1 "github.com/fission/fission/pkg/apis/core/v1" - "github.com/fission/fission/pkg/mqtrigger/factory" - "github.com/fission/fission/pkg/mqtrigger/messageQueue" - "github.com/fission/fission/pkg/mqtrigger/validator" - "github.com/fission/fission/pkg/utils" -) - -var natsClusterID string -var natsQueueGroup string -var natsClientID string - -func init() { - natsClusterID = os.Getenv("MESSAGE_QUEUE_CLUSTER_ID") - if natsClusterID == "" { - natsClusterID = defaultNatsClusterID - } - natsClientID = os.Getenv("MESSAGE_QUEUE_CLIENT_ID") - if natsClientID == "" { - natsClientID = defaultNatsClientID - } - natsQueueGroup = os.Getenv("MESSAGE_QUEUE_QUEUE_GROUP") - if natsQueueGroup == "" { - natsQueueGroup = defaultNatsQueueGroup - } - factory.Register(fv1.MessageQueueTypeNats, &Factory{}) - validator.Register(fv1.MessageQueueTypeNats, IsTopicValid) -} - -const ( - defaultNatsClusterID = "fissionMQTrigger" - defaultNatsClientID = "fission" - defaultNatsQueueGroup = "fission-messageQueueNatsTrigger" -) - -type ( - Nats struct { - logger *zap.Logger - nsConn ns.Conn - routerUrl string - } - - Factory struct{} -) - -func (factory *Factory) Create(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messageQueue.MessageQueue, error) { - return New(logger, mqCfg, routerUrl) -} - -func New(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messageQueue.MessageQueue, error) { - conn, err := ns.Connect(natsClusterID, natsClientID, ns.NatsURL(mqCfg.Url), - ns.SetConnectionLostHandler(func(conn ns.Conn, reason error) { - // TODO: Better way to handle connection lost problem. - // Currently, MessageQueue has no such interface to expose the status of underlying - // messaging service, hence MessageQueueTriggerManager has no way to detect and handle - // such situation properly. It takes some time to redesign interface of MessageQueue. - // For now, we simply fatal here. - logger.Fatal("Connection lost", zap.Error(reason)) - }), - ) - if err != nil { - return nil, err - } - nats := Nats{ - logger: logger.Named("nats"), - nsConn: conn, - routerUrl: routerUrl, - } - return nats, nil -} - -func (nats Nats) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Subscription, error) { - subj := trigger.Spec.Topic - - if !IsTopicValid(subj) { - return nil, fmt.Errorf("not a valid topic: %q", trigger.Spec.Topic) - } - - opts := []ns.SubscriptionOption{ - // Create a durable subscription to nats, so that triggers could retrieve last unack message. - // https://github.com/nats-io/stan.go#durable-subscriptions - ns.DurableName(string(trigger.ObjectMeta.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 messageQueue.Subscription) error { - return subscription.(ns.Subscription).Close() -} - -func msgHandler(nats *Nats, trigger *fv1.MessageQueueTrigger) func(*ns.Msg) { - return func(msg *ns.Msg) { - - // Support other function ref types - if trigger.Spec.FunctionReference.Type != fv1.FunctionReferenceTypeFunctionName { - nats.logger.Fatal("unsupported function reference type for trigger", - zap.Any("function_reference_type", trigger.Spec.FunctionReference.Type), - zap.String("trigger", trigger.ObjectMeta.Name)) - } - - // with the addition of multi-tenancy, the users can create functions in any namespace. however, - // the triggers can only be created in the same namespace as the function. - // so essentially, function namespace = trigger namespace. - url := nats.routerUrl + "/" + strings.TrimPrefix(utils.UrlForFunction(trigger.Spec.FunctionReference.Name, trigger.ObjectMeta.Namespace), "/") - nats.logger.Debug("making HTTP request", zap.String("url", url)) - - headers := map[string]string{ - "X-Fission-MQTrigger-Topic": trigger.Spec.Topic, - "X-Fission-MQTrigger-RespTopic": trigger.Spec.ResponseTopic, - "X-Fission-MQTrigger-ErrorTopic": trigger.Spec.ErrorTopic, - "Content-Type": trigger.Spec.ContentType, - } - - // Create request - req, err := http.NewRequest("POST", url, bytes.NewReader(msg.Data)) - - if err != nil { - nats.logger.Error("failed to create HTTP request to invoke function", - zap.Error(err), - zap.String("function_url", url)) - return - } - - for k, v := range headers { - req.Header.Set(k, v) - } - - var resp *http.Response - for attempt := 0; attempt <= trigger.Spec.MaxRetries; attempt++ { - // Make the request - resp, err = http.DefaultClient.Do(req) - if err != nil { - nats.logger.Error("sending function invocation request failed", - zap.Error(err), - zap.String("function_url", url), - zap.String("trigger", trigger.ObjectMeta.Name)) - continue - } - if resp == nil { - continue - } - if err == nil && resp.StatusCode == http.StatusOK { - // Success, quit retrying - break - } - } - - if resp == nil { - nats.logger.Warn("every function invocation retry failed; final retry gave empty response", - zap.String("function_url", url), - zap.String("trigger", trigger.ObjectMeta.Name)) - return - } - - defer resp.Body.Close() - - body, bodyErr := io.ReadAll(resp.Body) - if bodyErr != nil { - nats.logger.Error("error reading function invocation response", - zap.Error(err), - zap.String("function_url", url), - zap.String("trigger", trigger.ObjectMeta.Name)) - return - } - - // Only the latest error response will be published to error topic - if err != nil || resp.StatusCode != 200 { - if len(trigger.Spec.ErrorTopic) > 0 && len(body) > 0 { - publishErr := nats.nsConn.Publish(trigger.Spec.ErrorTopic, body) - if publishErr != nil { - nats.logger.Error("failed to publish function invocation error to error topic", - zap.Error(publishErr), - zap.String("topic", trigger.Spec.ErrorTopic), - zap.String("function_url", url), - zap.String("trigger", trigger.ObjectMeta.Name)) - // TODO: We will ack this message after max retries to prevent re-processing but - // this may cause message loss - } - } - return - } - - // Trigger acks message only if a request was processed successfully - err = msg.Ack() - if err != nil { - nats.logger.Error("failed to ack message after successful function invocation from trigger", - zap.Error(err), - zap.String("function_url", url), - zap.String("trigger", trigger.ObjectMeta.Name)) - } - - if len(trigger.Spec.ResponseTopic) > 0 { - err = nats.nsConn.Publish(trigger.Spec.ResponseTopic, body) - if err != nil { - nats.logger.Error("failed to publish message with function invocation response to topic", - zap.Error(err), - zap.String("topic", trigger.Spec.ResponseTopic), - zap.String("trigger", trigger.ObjectMeta.Name)) - } - } - } -} - -func IsTopicValid(topic string) bool { - // nats-streaming does not support wildcard channel. - return nsUtil.IsChannelNameValid(topic, false) -} diff --git a/skaffold.yaml b/skaffold.yaml index e86f31f5..b75f7786 100644 --- a/skaffold.yaml +++ b/skaffold.yaml @@ -44,7 +44,6 @@ deploy: namespace: fission pprof.enabled: false canaryDeployment.enabled: false - nats.enabled: false influxdb.enabled: false pruneInterval: "60" repository: index.docker.io diff --git a/test/kind_CI.sh b/test/kind_CI.sh index 7c6d132f..1853820d 100755 --- a/test/kind_CI.sh +++ b/test/kind_CI.sh @@ -81,8 +81,6 @@ main() { $ROOT/test/tests/test_specs/test_spec_archive/test_spec_archive.sh \ $ROOT/test/tests/test_environments/test_tensorflow_serving_env.sh \ $ROOT/test/tests/test_environments/test_go_env.sh \ - $ROOT/test/tests/mqtrigger/nats/test_mqtrigger.sh \ - $ROOT/test/tests/mqtrigger/nats/test_mqtrigger_error.sh \ $ROOT/test/tests/test_huge_response/test_huge_response.sh \ $ROOT/test/tests/test_kubectl/test_kubectl.sh \ $ROOT/test/tests/websocket/test_ws.sh @@ -126,4 +124,4 @@ main echo "Total Failures" $FAILURES if [[ $FAILURES != '0' ]]; then exit 1 -fi \ No newline at end of file +fi diff --git a/test/test_utils.sh b/test/test_utils.sh index b71930d3..da69b66b 100755 --- a/test/test_utils.sh +++ b/test/test_utils.sh @@ -203,7 +203,6 @@ set_environment() { # fission env export FISSION_URL=http://$(kubectl -n $ns get svc controller -o jsonpath='{...ip}') export FISSION_ROUTER=$(kubectl -n $ns get svc router -o jsonpath='{...ip}') - export FISSION_NATS_STREAMING_URL="http://defaultFissionAuthToken@$(kubectl -n $ns get svc nats-streaming -o jsonpath='{...ip}:{.spec.ports[0].port}')" # ingress controller env export INGRESS_CONTROLLER=$(kubectl -n ingress-nginx get svc ingress-nginx -o jsonpath='{...ip}') @@ -480,7 +479,6 @@ dump_logs() { dump_fission_logs $ns $fns executor dump_fission_logs $ns $fns storagesvc dump_fission_logs $ns $fns mqtrigger - dump_fission_logs $ns $fns mqtrigger-nats-streaming dump_function_pod_logs $ns $fns dump_builder_pod_logs $bns dump_fission_crds @@ -536,8 +534,6 @@ run_all_tests() { $ROOT/test/tests/test_specs/test_spec_archive/test_spec_archive.sh \ $ROOT/test/tests/test_environments/test_tensorflow_serving_env.sh \ $ROOT/test/tests/test_environments/test_go_env.sh \ - $ROOT/test/tests/mqtrigger/nats/test_mqtrigger.sh \ - $ROOT/test/tests/mqtrigger/nats/test_mqtrigger_error.sh \ $ROOT/test/tests/test_huge_response/test_huge_response.sh \ $ROOT/test/tests/test_kubectl/test_kubectl.sh $ROOT/test/tests/websocket/test_ws.sh diff --git a/test/tests/mqtrigger/nats/main.js b/test/tests/mqtrigger/nats/main.js deleted file mode 100644 index b791d1af..00000000 --- a/test/tests/mqtrigger/nats/main.js +++ /dev/null @@ -1,6 +0,0 @@ -module.exports = async function(context) { - return { - status: 200, - body: "Hello, World!" - }; -} \ No newline at end of file diff --git a/test/tests/mqtrigger/nats/main_error.js b/test/tests/mqtrigger/nats/main_error.js deleted file mode 100644 index 37466bc9..00000000 --- a/test/tests/mqtrigger/nats/main_error.js +++ /dev/null @@ -1,6 +0,0 @@ -module.exports = async function(context) { - return { - status: 400, - body: "Hello, World!" - }; -} \ No newline at end of file diff --git a/test/tests/mqtrigger/nats/stan-pub/main.go b/test/tests/mqtrigger/nats/stan-pub/main.go deleted file mode 100644 index 1463de9a..00000000 --- a/test/tests/mqtrigger/nats/stan-pub/main.go +++ /dev/null @@ -1,143 +0,0 @@ -// This file originally came from official Nats.io GitHub repository. -// You can reach the original file with the following link: -// https://github.com/nats-io/stan.go/blob/master/examples/stan-pub/main.go - -// Copyright 2016-2019 The NATS 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 main - -import ( - "flag" - "fmt" - "log" - "os" - "sync" - "time" - - nats "github.com/nats-io/nats.go" - "github.com/nats-io/stan.go" -) - -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 - -cr, --creds NATS 2.0 Credentials -` - -// 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 - clientID string - URL string - async bool - userCreds 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") - flag.StringVar(&userCreds, "cr", "", "Credentials File") - flag.StringVar(&userCreds, "creds", "", "Credentials File") - - log.SetFlags(0) - flag.Usage = usage - flag.Parse() - - args := flag.Args() - - if len(args) < 1 { - usage() - } - - // Connect Options. - opts := []nats.Option{nats.Name("NATS Streaming Example Publisher")} - // Use UserCredentials - if userCreds != "" { - opts = append(opts, nats.UserCredentials(userCreds)) - } - - // Connect to NATS - nc, err := nats.Connect(URL, opts...) - if err != nil { - log.Fatal(err) - } - defer nc.Close() - - sc, err := stan.Connect(clusterID, clientID, stan.NatsConn(nc)) - 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 { - 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/tests/mqtrigger/nats/stan-sub/main.go b/test/tests/mqtrigger/nats/stan-sub/main.go deleted file mode 100644 index 40e65d57..00000000 --- a/test/tests/mqtrigger/nats/stan-sub/main.go +++ /dev/null @@ -1,188 +0,0 @@ -// This file originally came from official Nats.io GitHub repository. -// You can reach the original file with the following link: -// https://github.com/nats-io/stan.go/blob/master/examples/stan-sub/main.go - -// Copyright 2016-2019 The NATS 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 main - -import ( - "flag" - "fmt" - "log" - "os" - "os/signal" - "syscall" - "time" - - nats "github.com/nats-io/nats.go" - "github.com/nats-io/stan.go" - "github.com/nats-io/stan.go/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 - -cr, --creds NATS 2.0 Credentials - -Subscription Options: - --qgroup Queue group - --all Deliver all available messages - --last Deliver starting with last published message - --since Deliver messages in last interval (e.g. 1s, 1hr) - --seq Start at seqno - --new_only Only deliver new messages - --durable Durable subscriber name - --unsub 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, i int) { - log.Printf("[#%d] Received: %s\n", i, m) -} - -func main() { - var ( - clusterID, clientID string - URL string - userCreds string - showTime bool - qgroup string - unsubscribe bool - startSeq uint64 - startDelta string - deliverAll bool - newOnly bool - deliverLast bool - durable 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-sub", "The NATS Streaming client ID to connect with") - flag.StringVar(&clientID, "clientid", "stan-sub", "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", true, "Deliver all") - flag.BoolVar(&newOnly, "new_only", false, "Only new messages") - 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, "unsub", false, "Unsubscribe the durable on exit") - flag.BoolVar(&unsubscribe, "unsubscribe", false, "Unsubscribe the durable on exit") - flag.StringVar(&userCreds, "cr", "", "Credentials File") - flag.StringVar(&userCreds, "creds", "", "Credentials File") - - log.SetFlags(0) - flag.Usage = usage - flag.Parse() - - args := flag.Args() - - if len(args) < 1 { - log.Printf("Error: A subject must be specified.") - usage() - } - - // Connect Options. - opts := []nats.Option{nats.Name("NATS Streaming Example Subscriber")} - // Use UserCredentials - if userCreds != "" { - opts = append(opts, nats.UserCredentials(userCreds)) - } - - // Connect to NATS - nc, err := nats.Connect(URL, opts...) - if err != nil { - log.Fatal(err) - } - defer nc.Close() - - sc, err := stan.Connect(clusterID, clientID, stan.NatsConn(nc), - stan.SetConnectionLostHandler(func(_ stan.Conn, reason error) { - log.Fatalf("Connection lost, reason: %v", reason) - })) - 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) - - // Process Subscriber Options. - startOpt := stan.StartAt(pb.StartPosition_NewOnly) - if startSeq != 0 { - startOpt = stan.StartAtSequence(startSeq) - } else if deliverLast { - startOpt = stan.StartWithLastReceived() - } else if deliverAll && !newOnly { - startOpt = stan.DeliverAllAvailable() - } else if startDelta != "" { - ago, err := time.ParseDuration(startDelta) - if err != nil { - sc.Close() - log.Fatal(err) - } - startOpt = stan.StartAtTimeDelta(ago) - } - - subj, i := args[0], 0 - mcb := func(msg *stan.Msg) { - i++ - printMsg(msg, i) - } - - sub, err := sc.QueueSubscribe(subj, qgroup, mcb, startOpt, stan.DurableName(durable)) - if err != nil { - sc.Close() - log.Fatal(err) - } - - log.Printf("Listening on [%s], clientID=[%s], qgroup=[%s] durable=[%s]\n", subj, clientID, qgroup, durable) - - if showTime { - log.SetFlags(log.LstdFlags) - } - - // Wait for a SIGINT (perhaps triggered by user with CTRL-C) - // Run cleanup when signal is received - signalChan := make(chan os.Signal, 1) - cleanupDone := make(chan bool) - signal.Notify(signalChan, syscall.SIGINT, syscall.SIGTERM) - go func() { - for sig := range signalChan { - fmt.Printf("\nReceived signal %s, unsubscribing and closing connection...\n\n", sig.String()) - // Do not unsubscribe a durable on exit, except if asked to. - if durable == "" || unsubscribe { - err := sub.Unsubscribe() - if err != nil { - log.Fatal(err) - } - } - sc.Close() - cleanupDone <- true - } - }() - <-cleanupDone -} diff --git a/test/tests/mqtrigger/nats/test_mqtrigger.sh b/test/tests/mqtrigger/nats/test_mqtrigger.sh deleted file mode 100755 index ce36b42a..00000000 --- a/test/tests/mqtrigger/nats/test_mqtrigger.sh +++ /dev/null @@ -1,72 +0,0 @@ -#!/bin/bash -#test:disabled - -# -# Create a function and trigger it using NATS -# - -set -euo pipefail -source $(dirname $0)/../../../utils.sh -set +x - -TEST_ID=$(generate_test_id) -echo "TEST_ID = $TEST_ID" - -ROOT=$(dirname $0)/../../.. -DIR=$(dirname $0) - -clusterID="fissionMQTrigger" -pubClientID="clientPub-$TEST_ID" -subClientID="clientSub-$TEST_ID" -topic="foo.bar$TEST_ID" -resptopic="foo.foo$TEST_ID" -#FISSION_NATS_STREAMING_URL="http://defaultFissionAuthToken@$(minikube ip):4222" -expectedRespOutput="subject:\"$resptopic\" data:\"Hello, World!\"" - -env=nodejs-$TEST_ID -fn=hello-$TEST_ID -mqt=mqt-$TEST_ID - -cleanup() { - log "Cleaning up..." - clean_resource_by_id $TEST_ID -} - -if [ -z "${TEST_NOCLEANUP:-}" ]; then - trap cleanup EXIT -else - log "TEST_NOCLEANUP is set; not cleaning up test artifacts afterwards." -fi - -log "Creating nodejs env" -fission env create --name $env --image $NODE_RUNTIME_IMAGE - -log "Creating function" -fission fn create --name $fn --env $env --code $DIR/main.js --method GET - -log "Creating message queue trigger" -fission mqtrigger create --name $mqt --function $fn --mqtype "nats-streaming" --topic $topic --resptopic $resptopic - -# wait until nats trigger is created -sleep 5 - -# -# Send a message -# -log "Sending message" -go run $DIR/stan-pub/main.go -s $FISSION_NATS_STREAMING_URL -c $clusterID -id $pubClientID $topic "" - -# -# Wait for message on response topic -# -log "Waiting for response" -response=$(timeout 10s go run $DIR/stan-sub/main.go --last -s $FISSION_NATS_STREAMING_URL -c $clusterID -id $subClientID $resptopic 2>&1 || true) -log "Output from subscriber" -echo "$response" -echo "$response" | grep "$expectedRespOutput" - -log "Deleting message queue trigger" -fission mqtrigger delete --name $mqt - -log "Subscriber received expected response: $response" -log "Test PASSED" diff --git a/test/tests/mqtrigger/nats/test_mqtrigger_error.sh b/test/tests/mqtrigger/nats/test_mqtrigger_error.sh deleted file mode 100755 index 8956fcfb..00000000 --- a/test/tests/mqtrigger/nats/test_mqtrigger_error.sh +++ /dev/null @@ -1,75 +0,0 @@ -#!/bin/bash -#test:disabled - -# -# Create a function and trigger it using NATS -# To run this on Minikube, uncomment line 24 - -set -euo pipefail -set +x -source $(dirname $0)/../../../utils.sh - -TEST_ID=$(generate_test_id) -echo "TEST_ID = $TEST_ID" - -ROOT=$(dirname $0)/../../.. -DIR=$(dirname $0) - -clusterID="fissionMQTrigger" -pubClientID="clientPub-$TEST_ID" -subClientID="clientSub-$TEST_ID" -topic="foo.bar$TEST_ID" -resptopic="foo.foo$TEST_ID" -errortopic="foo.error$TEST_ID" -maxretries=1 -#FISSION_NATS_STREAMING_URL="http://defaultFissionAuthToken@$(minikube ip):4222" -expectedRespOutput="subject:\"$errortopic\" data:\"Hello, World!\"" - -env=nodejs-$TEST_ID -fn=hello-$TEST_ID -mqt=mqt-$TEST_ID - -cleanup() { - log "Cleaning up..." - clean_resource_by_id $TEST_ID -} - -if [ -z "${TEST_NOCLEANUP:-}" ]; then - trap cleanup EXIT -else - log "TEST_NOCLEANUP is set; not cleaning up test artifacts afterwards." -fi - -log "Creating nodejs env" -fission env create --name $env --image $NODE_RUNTIME_IMAGE - -log "Creating function" -fission fn create --name $fn --env $env --code $DIR/main_error.js --method GET - -log "Creating message queue trigger" -fission mqtrigger create --name $mqt --function $fn --mqtype "nats-streaming" --topic $topic --resptopic $resptopic --errortopic $errortopic --maxretries $maxretries -log "Updated mqtrigger list" -fission mqtrigger list - -# wait until nats trigger is created -sleep 5 - -# -# Send a message -# -log "Sending message" -go run $DIR/stan-pub/main.go -s $FISSION_NATS_STREAMING_URL -c $clusterID -id $pubClientID $topic "" - -# -# Wait for message on error topic -# -log "Waiting for response" -response=$(timeout 10s go run $DIR/stan-sub/main.go --last -s $FISSION_NATS_STREAMING_URL -c $clusterID -id $subClientID $errortopic 2>&1 || true) -log "Output from subscriber" -echo "$response" -echo "$response" | grep "$expectedRespOutput" - -log "Deleting message queue trigger" -fission mqtrigger delete --name $mqt - -log "Test PASSED" diff --git a/tools/port-forward-nats.sh b/tools/port-forward-nats.sh deleted file mode 100755 index 30f0a332..00000000 --- a/tools/port-forward-nats.sh +++ /dev/null @@ -1,23 +0,0 @@ -#!/bin/bash - -namespace=$1 -if [ -z "$namespace" ] -then - namespace=fission -fi - -svc=$1 -if [ -z "$svc" ] -then - svc=nats-streaming -fi - -port=$2 -if [ -z "$port" ] -then - port=8888 -fi - -kubectl get pods -l svc=$svc -o name --namespace $namespace | \ - sed 's/^.*\///' | \ - xargs -I{} kubectl port-forward {} $port:$port -n $namespace &