diff --git a/console/deploy/console.yaml b/console/deploy/console.yaml index cf2be2c..491177f 100644 --- a/console/deploy/console.yaml +++ b/console/deploy/console.yaml @@ -15,9 +15,12 @@ rules: - apiGroups: [""] resources: ["pods/log"] verbs: ["get"] + - apiGroups: [""] + resources: ["secrets"] + verbs: ["get", "list", "create", "delete"] - apiGroups: ["apps"] resources: ["deployments"] - verbs: ["get", "list", "update", "patch"] + verbs: ["get", "list", "create", "update", "patch", "delete"] - apiGroups: ["fission.io"] resources: ["environments", "packages", "functions", "httptriggers", "timetriggers"] verbs: ["get", "list", "create", "update", "patch", "delete"] @@ -55,7 +58,7 @@ spec: serviceAccountName: fission-console containers: - name: console - image: naeel/fission-console:v1.3.87 + image: naeel/fission-console:v1.3.88 imagePullPolicy: Always ports: - containerPort: 8090 diff --git a/console/internal/api/kwtriggers.go b/console/internal/api/kwtriggers.go new file mode 100644 index 0000000..a33273c --- /dev/null +++ b/console/internal/api/kwtriggers.go @@ -0,0 +1,241 @@ +package api + +// kwtriggers.go — CRUD хендлеры для KubernetesWatchTrigger (Fission KW Trigger). +// +// РЕШЕНИЕ ПО АРХИТЕКТУРЕ (2026-05-11): +// KubernetesWatchTrigger позволяет вызывать функцию при изменении K8s объектов. +// spec.type — тип ресурса: Pod, Service, Deployment, ConfigMap, и т.д. +// spec.namespace — namespace для слежения (по умолчанию = namespace пользователя) +// spec.labelselector — label selector в формате "key=value,key2=value2" +// spec.functionref — ссылка на функцию +// +// ОСОБЕННОСТИ: +// - Fission kubewatcher компонент должен быть задеплоен. +// - namespace в spec — это WATCHED namespace (не namespace триггера). +// Для безопасности ограничиваем: только namespace пользователя или пустое (тогда = userNS). +// - labelselector опционален, "" = смотрим на все ресурсы типа resourceType в namespace. +// +// ОШИБКИ В ПРОЦЕССЕ: +// - (нет, первая реализация) + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "strings" + "time" + + "fission-console/internal/fission" + "fission-console/internal/model" + + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" +) + +// validKWResourceTypes — поддерживаемые типы ресурсов для слежения. +// Расширяемо — это не ограничение CRD, просто UI-валидация. +var validKWResourceTypes = map[string]struct{}{ + "pod": {}, + "service": {}, + "deployment": {}, + "configmap": {}, + "secret": {}, + "namespace": {}, + "replicaset": {}, + "statefulset": {}, + "daemonset": {}, + "job": {}, +} + +func (s *Server) handleKWTriggersRoot(w http.ResponseWriter, r *http.Request) { + switch r.Method { + case http.MethodGet: + s.handleList(fission.KWTrigGVR)(w, r) + case http.MethodPost: + s.handleCreateKWTrigger(w, r) + default: + http.Error(w, "method not allowed", http.StatusMethodNotAllowed) + } +} + +func (s *Server) handleKWTriggersAction(w http.ResponseWriter, r *http.Request) { + path := strings.TrimPrefix(r.URL.Path, "/console/api/kwtriggers/") + path = strings.TrimPrefix(path, "/api/kwtriggers/") + name := strings.Trim(path, "/") + if name == "" || strings.Contains(name, "/") { + http.NotFound(w, r) + return + } + + switch r.Method { + case http.MethodGet: + s.handleGetKWTrigger(w, r, name) + case http.MethodDelete: + s.handleDeleteKWTrigger(w, r, name) + default: + http.Error(w, "method not allowed", http.StatusMethodNotAllowed) + } +} + +func (s *Server) handleCreateKWTrigger(w http.ResponseWriter, r *http.Request) { + var req model.CreateKWTriggerRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("decode request: %v", err)) + return + } + + ns := s.userNS(r) + + // Если namespace не задан — используем namespace пользователя + if strings.TrimSpace(req.Namespace) == "" { + req.Namespace = ns + } + + // Безопасность: нельзя смотреть за чужим namespace + if req.Namespace != ns { + writeJSONError(w, http.StatusForbidden, "can only watch your own namespace") + return + } + + if err := validateKWTriggerRequest(req); err != nil { + writeJSONError(w, http.StatusBadRequest, err.Error()) + return + } + + ctx, cancel := context.WithTimeout(r.Context(), 15*time.Second) + defer cancel() + + // Проверяем что функция существует + if _, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Get(ctx, req.FunctionName, metav1.GetOptions{}); err != nil { + if apierrors.IsNotFound(err) { + writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("function %q not found", req.FunctionName)) + return + } + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get function: %v", err)) + return + } + + obj := buildKWTriggerObject(ns, req) + created, err := s.dyn.Resource(fission.KWTrigGVR).Namespace(ns).Create(ctx, obj, metav1.CreateOptions{}) + if err != nil { + if apierrors.IsAlreadyExists(err) { + writeJSONError(w, http.StatusConflict, fmt.Sprintf("kwtrigger %q already exists", req.Name)) + return + } + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create kwtrigger: %v", err)) + return + } + + writeAnyJSON(w, http.StatusCreated, kwTriggerResponse(created)) +} + +func (s *Server) handleGetKWTrigger(w http.ResponseWriter, r *http.Request, name string) { + ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) + defer cancel() + + obj, err := s.dyn.Resource(fission.KWTrigGVR).Namespace(s.userNS(r)).Get(ctx, name, metav1.GetOptions{}) + if err != nil { + status := http.StatusBadGateway + if apierrors.IsNotFound(err) { + status = http.StatusNotFound + } + writeJSONError(w, status, fmt.Sprintf("get kwtrigger %q: %v", name, err)) + return + } + + writeAnyJSON(w, http.StatusOK, kwTriggerResponse(obj)) +} + +func (s *Server) handleDeleteKWTrigger(w http.ResponseWriter, r *http.Request, name string) { + ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) + defer cancel() + + if err := s.dyn.Resource(fission.KWTrigGVR).Namespace(s.userNS(r)).Delete(ctx, name, metav1.DeleteOptions{}); err != nil { + status := http.StatusBadGateway + if apierrors.IsNotFound(err) { + status = http.StatusNotFound + } + writeJSONError(w, status, fmt.Sprintf("delete kwtrigger %q: %v", name, err)) + return + } + + writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name}) +} + +// --- Вспомогательные функции --- + +func validateKWTriggerRequest(req model.CreateKWTriggerRequest) error { + if strings.TrimSpace(req.Name) == "" { + return fmt.Errorf("name is required") + } + if strings.TrimSpace(req.FunctionName) == "" { + return fmt.Errorf("functionName is required") + } + rt := strings.ToLower(strings.TrimSpace(req.ResourceType)) + if _, ok := validKWResourceTypes[rt]; !ok { + return fmt.Errorf("resourceType must be one of: Pod, Service, Deployment, ConfigMap, Secret, Namespace, ReplicaSet, StatefulSet, DaemonSet, Job") + } + return nil +} + +func buildKWTriggerObject(ns string, req model.CreateKWTriggerRequest) *unstructured.Unstructured { + // Fission ожидает capitalize: Pod, Service, Deployment + resourceType := capitalize(strings.TrimSpace(req.ResourceType)) + + spec := map[string]any{ + "type": resourceType, + "namespace": req.Namespace, + "functionref": map[string]any{ + "type": "name", + "name": req.FunctionName, + }, + } + if req.LabelSelector != "" { + spec["labelselector"] = req.LabelSelector + } + + return &unstructured.Unstructured{Object: map[string]any{ + "apiVersion": "fission.io/v1", + "kind": "KubernetesWatchTrigger", + "metadata": map[string]any{ + "name": req.Name, + "namespace": ns, + }, + "spec": spec, + }} +} + +func kwTriggerResponse(obj *unstructured.Unstructured) map[string]any { + spec, _ := obj.Object["spec"].(map[string]any) + if spec == nil { + spec = map[string]any{} + } + fnref, _ := spec["functionref"].(map[string]any) + fnName := "" + if fnref != nil { + fnName, _ = fnref["name"].(string) + } + return map[string]any{ + "metadata": map[string]any{ + "name": obj.GetName(), + "namespace": obj.GetNamespace(), + }, + "spec": map[string]any{ + "resourceType": spec["type"], + "namespace": spec["namespace"], + "labelSelector": spec["labelselector"], + "functionName": fnName, + }, + } +} + +// capitalize приводит первый символ к верхнему регистру, остальное без изменений. +// "pod" → "Pod", "deployment" → "Deployment" +func capitalize(s string) string { + if s == "" { + return s + } + return strings.ToUpper(s[:1]) + strings.ToLower(s[1:]) +} diff --git a/console/internal/api/mqtriggers.go b/console/internal/api/mqtriggers.go new file mode 100644 index 0000000..cc6fef8 --- /dev/null +++ b/console/internal/api/mqtriggers.go @@ -0,0 +1,286 @@ +package api + +// mqtriggers.go — CRUD для MQ-триггеров через K8s Deployment + Secret. +// +// АРХИТЕКТУРА (2026-05-11): +// Вместо Fission MessageQueueTrigger CRD (требует mqtrigger компонент Kafka/NATS) +// Console деплоит собственный sqs-consumer Deployment в namespace пользователя. +// +// При CREATE: +// 1. Создаём K8s Secret (sqs-mq-) с SQS credentials +// 2. Создаём K8s Deployment (mq-) с образом naeel/sqs-consumer:v1.0 +// FUNCTION_URL = http://router.fission.svc.cluster.local/ +// Лейблы: app.kubernetes.io/managed-by=fission-console, component=mq-trigger +// +// При LIST: deployments -n -l component=mq-trigger +// При DELETE: удаляем Deployment + Secret +// +// ИЗМЕНЕНИЯ: +// v1 — использовал Fission MQ CRD (компонент отсутствует в кластере) +// v2 — K8s Deployment + наш sqs-consumer образ + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "strings" + "time" + + "fission-console/internal/model" + + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +const ( + sqsConsumerImage = "naeel/sqs-consumer:v1.0" + sqsDefaultEndpoint = "http://shared-sqs.shared-sqs.svc.cluster.local:4100" + fissionRouterBase = "http://router.fission.svc.cluster.local" + mqTriggerLabelKey = "component" + mqTriggerLabelVal = "mq-trigger" + mqManagedByLabel = "app.kubernetes.io/managed-by" + mqManagedByVal = "fission-console" +) + +func mqSecretName(name string) string { return "sqs-mq-" + name } +func mqDeployName(name string) string { return "mq-" + name } + +// ── HTTP хендлеры ──────────────────────────────────────────────────── + +func (s *Server) handleMQTriggersRoot(w http.ResponseWriter, r *http.Request) { + switch r.Method { + case http.MethodGet: + s.handleListMQTriggers(w, r) + case http.MethodPost: + s.handleCreateMQTrigger(w, r) + default: + http.Error(w, "method not allowed", http.StatusMethodNotAllowed) + } +} + +func (s *Server) handleMQTriggersAction(w http.ResponseWriter, r *http.Request) { + path := strings.TrimPrefix(r.URL.Path, "/console/api/mqtriggers/") + name := strings.Trim(path, "/") + if name == "" || strings.Contains(name, "/") { + http.NotFound(w, r) + return + } + switch r.Method { + case http.MethodDelete: + s.handleDeleteMQTrigger(w, r, name) + default: + http.Error(w, "method not allowed", http.StatusMethodNotAllowed) + } +} + +// ── LIST ──────────────────────────────────────────────────────────── + +func (s *Server) handleListMQTriggers(w http.ResponseWriter, r *http.Request) { + ns := s.userNS(r) + ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) + defer cancel() + + labelSel := fmt.Sprintf("%s=%s,%s=%s", mqManagedByLabel, mqManagedByVal, mqTriggerLabelKey, mqTriggerLabelVal) + deployList, err := s.kube.AppsV1().Deployments(ns).List(ctx, metav1.ListOptions{LabelSelector: labelSel}) + if err != nil { + writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("list mq deployments: %v", err)) + return + } + + items := make([]map[string]any, 0, len(deployList.Items)) + for i := range deployList.Items { + items = append(items, mqDeployToResponse(&deployList.Items[i])) + } + writeAnyJSON(w, http.StatusOK, map[string]any{"items": items}) +} + +// ── CREATE ────────────────────────────────────────────────────────── + +func (s *Server) handleCreateMQTrigger(w http.ResponseWriter, r *http.Request) { + var req model.CreateMQTriggerRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeJSONError(w, http.StatusBadRequest, "decode request: "+err.Error()) + return + } + if err := validateMQRequest(req); err != nil { + writeJSONError(w, http.StatusBadRequest, err.Error()) + return + } + + ns := s.userNS(r) + ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second) + defer cancel() + + endpoint := strings.TrimSpace(req.SqsEndpoint) + if endpoint == "" { + endpoint = sqsDefaultEndpoint + } + functionURL := fissionRouterBase + "/" + req.FunctionName + secretName := mqSecretName(req.Name) + deployName := mqDeployName(req.Name) + + // 1. Secret с SQS credentials + secret := &corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{ + Name: secretName, + Namespace: ns, + Labels: map[string]string{ + mqManagedByLabel: mqManagedByVal, + mqTriggerLabelKey: mqTriggerLabelVal, + "mq-trigger-name": req.Name, + }, + }, + StringData: map[string]string{ + "SQS_ACCESS_KEY": req.AccessKey, + "SQS_SECRET_KEY": req.SecretKey, + "SQS_ENDPOINT": endpoint, + }, + } + if _, err := s.kube.CoreV1().Secrets(ns).Create(ctx, secret, metav1.CreateOptions{}); err != nil { + if apierrors.IsAlreadyExists(err) { + writeJSONError(w, http.StatusConflict, fmt.Sprintf("mq trigger %q already exists", req.Name)) + return + } + writeJSONError(w, http.StatusBadGateway, "create secret: "+err.Error()) + return + } + + // 2. Deployment (sqs-consumer) + replicas := int32(1) + deploy := &appsv1.Deployment{ + ObjectMeta: metav1.ObjectMeta{ + Name: deployName, + Namespace: ns, + Labels: map[string]string{ + mqManagedByLabel: mqManagedByVal, + mqTriggerLabelKey: mqTriggerLabelVal, + "mq-trigger-name": req.Name, + }, + Annotations: map[string]string{ + "fission-console/mq-trigger-name": req.Name, + "fission-console/function": req.FunctionName, + "fission-console/queue": req.Queue, + "fission-console/sqs-endpoint": endpoint, + }, + }, + Spec: appsv1.DeploymentSpec{ + Replicas: &replicas, + Selector: &metav1.LabelSelector{ + MatchLabels: map[string]string{"mq-trigger-name": req.Name}, + }, + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{ + Labels: map[string]string{ + "mq-trigger-name": req.Name, + mqTriggerLabelKey: mqTriggerLabelVal, + }, + }, + Spec: corev1.PodSpec{ + Containers: []corev1.Container{{ + Name: "sqs-consumer", + Image: sqsConsumerImage, + ImagePullPolicy: corev1.PullAlways, + Env: []corev1.EnvVar{ + {Name: "SQS_QUEUE_NAME", Value: req.Queue}, + {Name: "SQS_REGION", Value: "us-east-1"}, + {Name: "FUNCTION_URL", Value: functionURL}, + {Name: "POLL_INTERVAL", Value: "5"}, + {Name: "MAX_MESSAGES", Value: "1"}, + {Name: "MAX_RETRIES", Value: "3"}, + }, + EnvFrom: []corev1.EnvFromSource{{ + SecretRef: &corev1.SecretEnvSource{ + LocalObjectReference: corev1.LocalObjectReference{Name: secretName}, + }, + }}, + Resources: corev1.ResourceRequirements{ + Limits: corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("50m"), + corev1.ResourceMemory: resource.MustParse("32Mi"), + }, + Requests: corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("10m"), + corev1.ResourceMemory: resource.MustParse("16Mi"), + }, + }, + }}, + }, + }, + }, + } + + created, err := s.kube.AppsV1().Deployments(ns).Create(ctx, deploy, metav1.CreateOptions{}) + if err != nil { + _ = s.kube.CoreV1().Secrets(ns).Delete(ctx, secretName, metav1.DeleteOptions{}) + writeJSONError(w, http.StatusBadGateway, "create deployment: "+err.Error()) + return + } + + writeAnyJSON(w, http.StatusCreated, mqDeployToResponse(created)) +} + +// ── DELETE ────────────────────────────────────────────────────────── + +func (s *Server) handleDeleteMQTrigger(w http.ResponseWriter, r *http.Request, name string) { + ns := s.userNS(r) + ctx, cancel := context.WithTimeout(r.Context(), 15*time.Second) + defer cancel() + + dErr := s.kube.AppsV1().Deployments(ns).Delete(ctx, mqDeployName(name), metav1.DeleteOptions{}) + sErr := s.kube.CoreV1().Secrets(ns).Delete(ctx, mqSecretName(name), metav1.DeleteOptions{}) + + if dErr != nil && !apierrors.IsNotFound(dErr) { + writeJSONError(w, http.StatusBadGateway, "delete deployment: "+dErr.Error()) + return + } + if sErr != nil && !apierrors.IsNotFound(sErr) { + writeJSONError(w, http.StatusBadGateway, "delete secret: "+sErr.Error()) + return + } + writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name}) +} + +// ── Вспомогательные ───────────────────────────────────────────────── + +func validateMQRequest(req model.CreateMQTriggerRequest) error { + if strings.TrimSpace(req.Name) == "" { + return fmt.Errorf("name is required") + } + if strings.TrimSpace(req.FunctionName) == "" { + return fmt.Errorf("functionName is required") + } + if strings.TrimSpace(req.Queue) == "" { + return fmt.Errorf("queue is required") + } + if strings.TrimSpace(req.AccessKey) == "" { + return fmt.Errorf("accessKey is required") + } + if strings.TrimSpace(req.SecretKey) == "" { + return fmt.Errorf("secretKey is required") + } + return nil +} + +func mqDeployToResponse(d *appsv1.Deployment) map[string]any { + ann := d.Annotations + if ann == nil { + ann = map[string]string{} + } + triggerName := ann["fission-console/mq-trigger-name"] + if triggerName == "" { + triggerName = strings.TrimPrefix(d.Name, "mq-") + } + return map[string]any{ + "name": triggerName, + "deployName": d.Name, + "functionName": ann["fission-console/function"], + "queue": ann["fission-console/queue"], + "sqsEndpoint": ann["fission-console/sqs-endpoint"], + "ready": d.Status.ReadyReplicas > 0, + "replicas": d.Status.ReadyReplicas, + } +} diff --git a/console/internal/api/server.go b/console/internal/api/server.go index 09f0e2c..91625db 100644 --- a/console/internal/api/server.go +++ b/console/internal/api/server.go @@ -169,6 +169,10 @@ func (s *Server) RegisterRoutes(mux *http.ServeMux) { mux.HandleFunc("/console/api/httptriggers", auth(s.handleList(fission.HTTPTrigGVR))) mux.HandleFunc("/console/api/timetriggers", auth(s.handleTimeTriggersRoot)) mux.HandleFunc("/console/api/timetriggers/", auth(s.handleTimeTriggersAction)) + mux.HandleFunc("/console/api/mqtriggers", auth(s.handleMQTriggersRoot)) + mux.HandleFunc("/console/api/mqtriggers/", auth(s.handleMQTriggersAction)) + mux.HandleFunc("/console/api/kwtriggers", auth(s.handleKWTriggersRoot)) + mux.HandleFunc("/console/api/kwtriggers/", auth(s.handleKWTriggersAction)) mux.HandleFunc("/console/api/ns/status", auth(s.handleNSStatus)) mux.HandleFunc("/console/api/ns/debug", auth(s.handleNSDebug)) mux.HandleFunc("/console/api/stats/dashboard-url", auth(s.handleStatsDashboard)) diff --git a/console/internal/fission/client.go b/console/internal/fission/client.go index 388b290..366038e 100644 --- a/console/internal/fission/client.go +++ b/console/internal/fission/client.go @@ -12,6 +12,8 @@ var ( FunctionGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "functions"} HTTPTrigGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "httptriggers"} TimeTrigGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "timetriggers"} + MQTrigGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "messagequeuetriggers"} + KWTrigGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "kuberneteswatchtriggers"} NamespaceGVR = schema.GroupVersionResource{Group: "", Version: "v1", Resource: "namespaces"} DeploymentGVR = schema.GroupVersionResource{Group: "apps", Version: "v1", Resource: "deployments"} diff --git a/console/internal/model/types.go b/console/internal/model/types.go index 1522a15..3f3ffd9 100644 --- a/console/internal/model/types.go +++ b/console/internal/model/types.go @@ -26,6 +26,28 @@ type CreateTimeTriggerRequest struct { SubPath string `json:"subpath"` } +// CreateMQTriggerRequest — тело POST /console/api/mqtriggers. +// Архитектура: Console создаёт K8s Deployment (sqs-consumer) + Secret с credentials +// в namespace пользователя. sqs-consumer поллит shared-sqs → вызывает Fission-функцию. +type CreateMQTriggerRequest struct { + Name string `json:"name"` + FunctionName string `json:"functionName"` + Queue string `json:"queue"` // имя очереди в SQS + SqsEndpoint string `json:"sqsEndpoint"` // URL SQS сервиса (default: internal shared-sqs) + AccessKey string `json:"accessKey"` // SQS access key тенанта + SecretKey string `json:"secretKey"` // SQS secret key тенанта +} + +// CreateKWTriggerRequest — тело POST /console/api/kwtriggers. +// Документация полей: kubectl get crd kuberneteswatchtriggers.fission.io -o json +type CreateKWTriggerRequest struct { + Name string `json:"name"` + FunctionName string `json:"functionName"` + ResourceType string `json:"resourceType"` // Pod, Service, Deployment и т.д. + Namespace string `json:"namespace"` // пустое = namespace пользователя + LabelSelector string `json:"labelSelector"` // "app=foo" или "" для всех +} + // UpdateCodeRequest — тело PUT /console/api/functions/:name/code. type UpdateCodeRequest struct { Code string `json:"code"` diff --git a/console/ui/index.html b/console/ui/index.html index e98a19a..905eed1 100644 --- a/console/ui/index.html +++ b/console/ui/index.html @@ -25,6 +25,7 @@ + @@ -102,7 +103,7 @@
NUBES
FISSION CONSOLE
-
v1.3.87
+
v1.3.88
@@ -136,6 +137,10 @@
Крон-функции
-
+
+
MQ-триггеры
+
-
+
@@ -174,6 +179,72 @@
Изменения применяются напрямую через CRD Fission.
+ + +
+
+
MQ-триггеры
+ +
+ + + + + + + + + + + + + + +
ИмяОчередьФункцияEndpoint SQSСтатусДействия
Загрузка...
+
MQ-триггер поллит очередь SQS и вызывает Fission-функцию при появлении сообщений.
+
+ + + + @@ -543,7 +614,7 @@
- v1.3.87 + v1.3.88
diff --git a/console/ui/js/app.js b/console/ui/js/app.js index 3c7a855..df4c2b4 100644 --- a/console/ui/js/app.js +++ b/console/ui/js/app.js @@ -16,6 +16,9 @@ async function reloadAll() { getJSON(API_BASE + '/timetriggers') ]); + // MQ-триггеры загружаем параллельно, не блокируем основную таблицу + if (typeof loadMQTriggers === 'function') loadMQTriggers(); + S.envs = envs || []; S.fns = fns || []; S.httpTriggers = http || []; diff --git a/console/ui/js/mq.js b/console/ui/js/mq.js new file mode 100644 index 0000000..b3a75ab --- /dev/null +++ b/console/ui/js/mq.js @@ -0,0 +1,103 @@ +// mq.js — MQ-триггеры (SQS → Fission function) +// Архитектура: Console создаёт K8s Deployment + Secret (sqs-consumer) в namespace пользователя. +// sqs-consumer поллит SQS очередь → HTTP POST в Fission-функцию → DeleteMessage. +// Backend: POST /console/api/mqtriggers, DELETE /console/api/mqtriggers/{name} +// v1.3.88 + +// ── Загрузка и отрисовка ──────────────────────────────────────────── + +async function loadMQTriggers() { + try { + const data = await apiFetch('/console/api/mqtriggers'); + renderMQTable(data.items || []); + const cnt = document.getElementById('mq-count'); + if (cnt) cnt.textContent = (data.items || []).length; + } catch (e) { + renderMQTable([]); + const cnt = document.getElementById('mq-count'); + if (cnt) cnt.textContent = '0'; + } +} + +function renderMQTable(items) { + const tbody = document.getElementById('mq-rows'); + if (!tbody) return; + if (!items.length) { + tbody.innerHTML = 'Нет MQ-триггеров'; + return; + } + tbody.innerHTML = items.map(t => ` + + ${escHtml(t.name)} + ${escHtml(t.queue)} + ${escHtml(t.functionName)} + ${escHtml(t.sqsEndpoint || '')} + ${t.ready ? '▶ Running' : '◼ Pending'} + + + + + `).join(''); +} + +// ── Создание ───────────────────────────────────────────────────────── + +function openCreateMQ() { + document.getElementById('mq-name').value = ''; + document.getElementById('mq-fn').value = ''; + document.getElementById('mq-queue').value = ''; + document.getElementById('mq-endpoint').value = 'http://shared-sqs.shared-sqs.svc.cluster.local:4100'; + document.getElementById('mq-access-key').value = ''; + document.getElementById('mq-secret-key').value = ''; + document.getElementById('mq-create-error').textContent = ''; + document.getElementById('mq-create-modal').style.display = 'flex'; +} + +function closeMQCreate() { + document.getElementById('mq-create-modal').style.display = 'none'; +} + +async function submitCreateMQ() { + const name = document.getElementById('mq-name').value.trim(); + const fnName = document.getElementById('mq-fn').value.trim(); + const queue = document.getElementById('mq-queue').value.trim(); + const endpoint = document.getElementById('mq-endpoint').value.trim(); + const accessKey = document.getElementById('mq-access-key').value.trim(); + const secretKey = document.getElementById('mq-secret-key').value.trim(); + const errEl = document.getElementById('mq-create-error'); + + if (!name || !fnName || !queue || !accessKey || !secretKey) { + errEl.textContent = 'Заполните все поля'; + return; + } + + errEl.textContent = ''; + try { + await apiFetch('/console/api/mqtriggers', { + method: 'POST', + body: JSON.stringify({ name, functionName: fnName, queue, sqsEndpoint: endpoint, accessKey, secretKey }) + }); + closeMQCreate(); + await loadMQTriggers(); + } catch (e) { + errEl.textContent = e.message || 'Ошибка создания'; + } +} + +// ── Удаление ───────────────────────────────────────────────────────── + +async function deleteMQTrigger(name) { + if (!confirm(`Удалить MQ-триггер "${name}"?`)) return; + try { + await apiFetch(`/console/api/mqtriggers/${encodeURIComponent(name)}`, { method: 'DELETE' }); + await loadMQTriggers(); + } catch (e) { + alert('Ошибка удаления: ' + (e.message || e)); + } +} + +// ── Утилита ────────────────────────────────────────────────────────── + +function escHtml(s) { + return String(s).replace(/&/g,'&').replace(//g,'>').replace(/"/g,'"'); +} diff --git a/examples/weather-demo/consumer/main.py b/examples/weather-demo/consumer/main.py index 56de1fe..9f58b7e 100644 --- a/examples/weather-demo/consumer/main.py +++ b/examples/weather-demo/consumer/main.py @@ -4,16 +4,22 @@ import psycopg2 def main(event, context): - pg_dsn = os.environ["PG_DSN"] + pg_dsn = os.environ.get( + "PG_DSN", + "postgresql://super:BQUF5ruECa1ZFlq4wYt3gPJUEmtBMkA9QNK4MM5Sd8al4ArMDlmT16DIKHYBPyif" + "@postgresqlk8s-master.dc5db45d-f8b4-4fd0-ad33-ec4dd017f2d5.svc.cluster.local:5432" + "/sqsdb?sslmode=disable", + ) # Извлекаем тело сообщения из SQS (POST от sqs-consumer) - body = getattr(event, "body", event) - if isinstance(body, (bytes, bytearray)): - body = body.decode("utf-8") - if isinstance(body, str): - data = json.loads(body) - else: - data = body + # event — Flask Request object: используем .data (bytes) или .get_json() + try: + data = event.get_json(force=True, silent=False) + except Exception: + raw = getattr(event, "data", None) or getattr(event, "body", b"") + if isinstance(raw, (bytes, bytearray)): + raw = raw.decode("utf-8") + data = json.loads(raw) if raw else {} conn = psycopg2.connect(pg_dsn) try: diff --git a/examples/weather-demo/fetcher/main.py b/examples/weather-demo/fetcher/main.py index a2affc3..64b1b07 100644 --- a/examples/weather-demo/fetcher/main.py +++ b/examples/weather-demo/fetcher/main.py @@ -1,26 +1,68 @@ import os import json +import time import requests import boto3 from botocore.config import Config - +# Open-Meteo: бесплатный API без ключа, реальные данные. +# https://open-meteo.com/en/docs CITIES = [ - "Moscow,RU", - "London,GB", - "Paphos,CY", - "Ulyanovsk,RU", - "Santiago,CL", + {"name": "Moscow", "country": "RU", "lat": 55.7558, "lon": 37.6173}, + {"name": "London", "country": "GB", "lat": 51.5074, "lon": -0.1278}, + {"name": "Paphos", "country": "CY", "lat": 34.7753, "lon": 32.4242}, + {"name": "Ulyanovsk", "country": "RU", "lat": 54.3282, "lon": 48.3866}, + {"name": "Santiago", "country": "CL", "lat": -33.4489, "lon": -70.6693}, ] +# WMO weather code → описание +WMO_DESCRIPTIONS = { + 0: "clear sky", 1: "mainly clear", 2: "partly cloudy", 3: "overcast", + 45: "fog", 48: "icy fog", 51: "light drizzle", 53: "drizzle", + 55: "heavy drizzle", 61: "light rain", 63: "rain", 65: "heavy rain", + 71: "light snow", 73: "snow", 75: "heavy snow", 80: "rain showers", + 81: "showers", 82: "violent showers", 95: "thunderstorm", +} + + +def fetch_city(city): + params = { + "latitude": city["lat"], + "longitude": city["lon"], + "current": "temperature_2m,apparent_temperature,relative_humidity_2m,surface_pressure,wind_speed_10m,weather_code", + "wind_speed_unit": "ms", + "timezone": "UTC", + } + resp = requests.get( + "https://api.open-meteo.com/v1/forecast", + params=params, + timeout=10, + ) + resp.raise_for_status() + cur = resp.json()["current"] + code = cur.get("weather_code", 0) + return { + "city": city["name"], + "country": city["country"], + "temperature": round(cur["temperature_2m"], 1), + "feels_like": round(cur["apparent_temperature"], 1), + "humidity": int(cur["relative_humidity_2m"]), + "pressure": int(cur["surface_pressure"]), + "wind_speed": round(cur["wind_speed_10m"], 1), + "description": WMO_DESCRIPTIONS.get(code, f"wmo:{code}"), + "owm_timestamp": int(time.time()), + } + def main(event, context): - api_key = os.environ["OWM_API_KEY"] sqs_endpoint = os.environ.get( "SQS_ENDPOINT", "http://shared-sqs.shared-sqs.svc.cluster.local:4100" ) - access_key = os.environ["SQS_ACCESS_KEY"] - secret_key = os.environ["SQS_SECRET_KEY"] + access_key = os.environ.get("SQS_ACCESS_KEY", "SSAK-a9964f2723bc6d347f48d153") + secret_key = os.environ.get( + "SQS_SECRET_KEY", + "2069e1ce05aaf94efe07aee18697352879e7626df239a9c71af0e9650b43bdd6", + ) queue_name = os.environ.get("SQS_QUEUE_NAME", "weather-data") cfg = Config(signature_version="s3v4", s3={"addressing_style": "path"}) @@ -33,7 +75,6 @@ def main(event, context): config=cfg, ) - # Создаём очередь если не существует try: queue_url = sqs.get_queue_url(QueueName=queue_name)["QueueUrl"] except Exception: @@ -42,27 +83,10 @@ def main(event, context): results = [] for city in CITIES: try: - resp = requests.get( - "https://api.openweathermap.org/data/2.5/weather", - params={"q": city, "appid": api_key, "units": "metric"}, - timeout=10, - ) - resp.raise_for_status() - data = resp.json() - msg = { - "city": data["name"], - "country": data["sys"]["country"], - "temperature": data["main"]["temp"], - "feels_like": data["main"]["feels_like"], - "humidity": data["main"]["humidity"], - "pressure": data["main"]["pressure"], - "wind_speed": data["wind"]["speed"], - "description": data["weather"][0]["description"], - "owm_timestamp": data["dt"], - } + msg = fetch_city(city) sqs.send_message(QueueUrl=queue_url, MessageBody=json.dumps(msg)) results.append({"city": msg["city"], "temp": msg["temperature"]}) except Exception as e: - results.append({"city": city, "error": str(e)}) + results.append({"city": city["name"], "error": str(e)}) return {"status": "ok", "sent": len(results), "results": results} diff --git a/examples/weather-demo/main.tf b/examples/weather-demo/main.tf index 1aab697..ab40183 100644 --- a/examples/weather-demo/main.tf +++ b/examples/weather-demo/main.tf @@ -24,11 +24,6 @@ variable "namespace" { description = "Namespace для функций и триггеров." } -variable "owm_api_key" { - sensitive = true - description = "OpenWeatherMap API key." -} - variable "sqs_access_key" { sensitive = true description = "SQS Access Key." @@ -61,7 +56,7 @@ resource "fission_iot_device" "weather_station" { resource "fission_environment" "python" { name = "weather-python" - image = "naeel/fission-python-env:v1.0" + image = "naeel/fission-python-env:v1.1" version = 2 namespace = var.namespace } @@ -73,7 +68,7 @@ resource "fission_package" "fetcher" { environment = fission_environment.python.name namespace = var.namespace source_dir = "${path.module}/fetcher" - deploy_type = "source" + deploy_type = "literal" } # ── Пакет: consumer ────────────────────────────────────────────────── @@ -83,20 +78,19 @@ resource "fission_package" "consumer" { environment = fission_environment.python.name namespace = var.namespace source_dir = "${path.module}/consumer" - deploy_type = "source" + deploy_type = "literal" } -# ── Функция: fetcher (читает OWM, пишет в SQS) ─────────────────────── - +# ── Функция: fetcher (читает Open-Meteo, пишет в SQS) ─────────────── +# Open-Meteo: бесплатный API без ключа. https://open-meteo.com resource "fission_function" "fetcher" { name = "weather-fetcher" environment = fission_environment.python.name namespace = var.namespace package_name = fission_package.fetcher.name entrypoint = "main" - # Env vars с секретами задаются через K8s Secret вне Terraform. - # Имя секрета: weather-fetcher-env (namespace: var.namespace) - # Ключи: OWM_API_KEY, SQS_ACCESS_KEY, SQS_SECRET_KEY + # Env vars задаются через K8s Secret weather-fetcher-env (namespace: var.namespace) + # Ключи: SQS_ACCESS_KEY, SQS_SECRET_KEY } # ── CRON trigger: каждые 10 минут ──────────────────────────────────── diff --git a/python-env/Dockerfile b/python-env/Dockerfile index a540dcd..5fafe3e 100644 --- a/python-env/Dockerfile +++ b/python-env/Dockerfile @@ -1,2 +1,3 @@ FROM ghcr.io/fission/python-env:latest +RUN pip install --no-cache-dir boto3==1.34.0 requests==2.31.0 psycopg2-binary==2.9.9 COPY server.py /app/server.py diff --git a/sqs-consumer/Dockerfile b/sqs-consumer/Dockerfile new file mode 100644 index 0000000..ed8fdf3 --- /dev/null +++ b/sqs-consumer/Dockerfile @@ -0,0 +1,16 @@ +# Dockerfile для sqs-consumer +# Назначение: поллит SQS-совместимую очередь (shared-sqs) и вызывает Fission-функцию. +# Образ: naeel/sqs-consumer:v1.0 +# 2026-05-11 + +FROM golang:1.22-alpine AS builder +WORKDIR /app +COPY go.mod ./ +COPY main.go ./ +RUN go mod tidy && CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -a -o sqs-consumer . + +FROM gcr.io/distroless/static:nonroot +WORKDIR / +COPY --from=builder /app/sqs-consumer . +USER 65532:65532 +ENTRYPOINT ["/sqs-consumer"] diff --git a/sqs-consumer/go.mod b/sqs-consumer/go.mod new file mode 100644 index 0000000..7970f42 --- /dev/null +++ b/sqs-consumer/go.mod @@ -0,0 +1,10 @@ +module sqs-consumer + +go 1.22 + +require ( + github.com/aws/aws-sdk-go-v2 v1.26.1 + github.com/aws/aws-sdk-go-v2/config v1.27.11 + github.com/aws/aws-sdk-go-v2/credentials v1.17.11 + github.com/aws/aws-sdk-go-v2/service/sqs v1.31.4 +) diff --git a/sqs-consumer/go.sum b/sqs-consumer/go.sum new file mode 100644 index 0000000..2ee0140 --- /dev/null +++ b/sqs-consumer/go.sum @@ -0,0 +1,28 @@ +github.com/aws/aws-sdk-go-v2 v1.26.1 h1:5554eUqIYVWpU0YmeeYZ0wU64H2VLBs8TlhRB2L+EkA= +github.com/aws/aws-sdk-go-v2 v1.26.1/go.mod h1:ffIFB97e2yNsv4aTSGkqtHnppsIJzw7G7BReUZ3jCXM= +github.com/aws/aws-sdk-go-v2/config v1.27.11 h1:f47rANd2LQEYHda2ddSCKYId18/8BhSRM4BULGmfgNA= +github.com/aws/aws-sdk-go-v2/config v1.27.11/go.mod h1:SMsV78RIOYdve1vf36z8LmnszlRWkwMQtomCAI0/mIE= +github.com/aws/aws-sdk-go-v2/credentials v1.17.11 h1:YuIB1dJNf1Re822rriUOTxopaHHvIq0l/pX3fwO+Tzs= +github.com/aws/aws-sdk-go-v2/credentials v1.17.11/go.mod h1:AQtFPsDH9bI2O+71anW6EKL+NcD7LG3dpKGMV4SShgo= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.1 h1:FVJ0r5XTHSmIHJV6KuDmdYhEpvlHpiSd38RQWhut5J4= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.1/go.mod h1:zusuAeqezXzAB24LGuzuekqMAEgWkVYukBec3kr3jUg= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.5 h1:aw39xVGeRWlWx9EzGVnhOR4yOjQDHPQ6o6NmBlscyQg= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.5/go.mod h1:FSaRudD0dXiMPK2UjknVwwTYyZMRsHv3TtkabsZih5I= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.5 h1:PG1F3OD1szkuQPzDw3CIQsRIrtTlUC3lP84taWzHlq0= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.5/go.mod h1:jU1li6RFryMz+so64PpKtudI+QzbKoIEivqdf6LNpOc= +github.com/aws/aws-sdk-go-v2/internal/ini v1.8.0 h1:hT8rVHwugYE2lEfdFE0QWVo81lF7jMrYJVDWI+f+VxU= +github.com/aws/aws-sdk-go-v2/internal/ini v1.8.0/go.mod h1:8tu/lYfQfFe6IGnaOdrpVgEL2IrrDOf6/m9RQum4NkY= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.11.2 h1:Ji0DY1xUsUr3I8cHps0G+XM3WWU16lP6yG8qu1GAZAs= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.11.2/go.mod h1:5CsjAbs3NlGQyZNFACh+zztPDI7fU6eW9QsxjfnuBKg= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.11.7 h1:ogRAwT1/gxJBcSWDMZlgyFUM962F51A5CRhDLbxLdmo= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.11.7/go.mod h1:YCsIZhXfRPLFFCl5xxY+1T9RKzOKjCut+28JSX2DnAk= +github.com/aws/aws-sdk-go-v2/service/sqs v1.31.4 h1:mE2ysZMEeQ3ulHWs4mmc4fZEhOfeY1o6QXAfDqjbSgw= +github.com/aws/aws-sdk-go-v2/service/sqs v1.31.4/go.mod h1:lCN2yKnj+Sp9F6UzpoPPTir+tSaC9Jwf6LcmTqnXFZw= +github.com/aws/aws-sdk-go-v2/service/sso v1.20.5 h1:vN8hEbpRnL7+Hopy9dzmRle1xmDc7o8tmY0klsr175w= +github.com/aws/aws-sdk-go-v2/service/sso v1.20.5/go.mod h1:qGzynb/msuZIE8I75DVRCUXw3o3ZyBmUvMwQ2t/BrGM= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.23.4 h1:Jux+gDDyi1Lruk+KHF91tK2KCuY61kzoCpvtvJJBtOE= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.23.4/go.mod h1:mUYPBhaF2lGiukDEjJX2BLRRKTmoUSitGDUgM4tRxak= +github.com/aws/aws-sdk-go-v2/service/sts v1.28.6 h1:cwIxeBttqPN3qkaAjcEcsh8NYr8n2HZPkcKgPAi1phU= +github.com/aws/aws-sdk-go-v2/service/sts v1.28.6/go.mod h1:FZf1/nKNEkHdGGJP/cI2MoIMquumuRK6ol3QQJNDxmw= +github.com/aws/smithy-go v1.20.2 h1:tbp628ireGtzcHDDmLT/6ADHidqnwgF57XOXZe6tp4Q= +github.com/aws/smithy-go v1.20.2/go.mod h1:krry+ya/rV9RDcV/Q16kpu6ypI4K2czasz0NC3qS14E= diff --git a/sqs-consumer/main.go b/sqs-consumer/main.go new file mode 100644 index 0000000..f12a8f3 --- /dev/null +++ b/sqs-consumer/main.go @@ -0,0 +1,256 @@ +// sqs-consumer — SQS poller → Fission function invoker +// +// Env vars: +// SQS_ENDPOINT — URL SQS сервиса +// SQS_QUEUE_NAME — имя очереди (обязательно) +// SQS_ACCESS_KEY — AccessKeyId тенанта +// SQS_SECRET_KEY — SecretAccessKey тенанта +// SQS_REGION — регион (default: us-east-1) +// FUNCTION_URL — полный URL функции +// ROUTER_USERNAME — логин для /auth/login роутера (optional) +// ROUTER_PASSWORD — пароль для /auth/login роутера (optional) +// ROUTER_LOGIN_URL — URL /auth/login (default: выводится из FUNCTION_URL) +// AUTH_TOKEN — статичный Bearer токен (если ROUTER_USERNAME не задан) +// POLL_INTERVAL — интервал поллинга в секундах (default: 5) +// MAX_MESSAGES — макс. сообщений за раз (default: 1) +// MAX_RETRIES — попыток вызова функции перед skip (default: 3) +// +// 2026-05-12 + +package main + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "log/slog" + "net/http" + "os" + "strconv" + "strings" + "sync" + "time" + + "github.com/aws/aws-sdk-go-v2/aws" + awsconfig "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/credentials" + "github.com/aws/aws-sdk-go-v2/service/sqs" +) + +// tokenManager управляет JWT токеном с автообновлением. +type tokenManager struct { + mu sync.Mutex + client *http.Client + loginURL string + username string + password string + staticToken string + cached string + expiresAt time.Time +} + +func newTokenManager(client *http.Client, functionURL, loginURL, username, password, staticToken string) *tokenManager { + if loginURL == "" && username != "" { + if i := strings.Index(functionURL, "://"); i >= 0 { + rest := functionURL[i+3:] + if j := strings.Index(rest, "/"); j >= 0 { + loginURL = functionURL[:i+3] + rest[:j] + "/auth/login" + } else { + loginURL = functionURL + "/auth/login" + } + } + } + return &tokenManager{ + client: client, + loginURL: loginURL, + username: username, + password: password, + staticToken: staticToken, + } +} + +func (tm *tokenManager) getToken(log *slog.Logger) string { + tm.mu.Lock() + defer tm.mu.Unlock() + + if tm.username == "" { + return tm.staticToken + } + if tm.cached != "" && time.Now().Add(30*time.Second).Before(tm.expiresAt) { + return tm.cached + } + token, exp := tm.login(log) + if token != "" { + tm.cached = token + tm.expiresAt = exp + log.Info("router token refreshed", "expiresAt", exp) + } + return tm.cached +} + +func (tm *tokenManager) login(log *slog.Logger) (string, time.Time) { + body, _ := json.Marshal(map[string]string{"username": tm.username, "password": tm.password}) + resp, err := tm.client.Post(tm.loginURL, "application/json", bytes.NewReader(body)) + if err != nil { + log.Error("router login failed", "err", err) + return "", time.Time{} + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated { + log.Error("router login non-2xx", "status", resp.StatusCode) + return "", time.Time{} + } + var result struct { + AccessToken string `json:"accesstoken"` + } + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil || result.AccessToken == "" { + log.Error("router login decode failed", "err", err) + return "", time.Time{} + } + return result.AccessToken, time.Now().Add(5 * time.Minute) +} + +func (tm *tokenManager) invalidate() { + tm.mu.Lock() + defer tm.mu.Unlock() + tm.cached = "" +} + +func main() { + log := slog.New(slog.NewJSONHandler(os.Stdout, nil)) + + endpoint := getenv("SQS_ENDPOINT", "http://shared-sqs.shared-sqs.svc.cluster.local:4100") + queueName := mustenv("SQS_QUEUE_NAME") + accessKey := mustenv("SQS_ACCESS_KEY") + secretKey := mustenv("SQS_SECRET_KEY") + region := getenv("SQS_REGION", "us-east-1") + functionURL := mustenv("FUNCTION_URL") + routerUser := getenv("ROUTER_USERNAME", "") + routerPass := getenv("ROUTER_PASSWORD", "") + routerLoginURL := getenv("ROUTER_LOGIN_URL", "") + staticToken := getenv("AUTH_TOKEN", "") + pollSec := parseInt(getenv("POLL_INTERVAL", "5"), 5) + maxMsg := int32(parseInt(getenv("MAX_MESSAGES", "1"), 1)) + maxRetries := parseInt(getenv("MAX_RETRIES", "3"), 3) + + customResolver := aws.EndpointResolverWithOptionsFunc( + func(service, reg string, opts ...interface{}) (aws.Endpoint, error) { + return aws.Endpoint{URL: endpoint, HostnameImmutable: true}, nil + }, + ) + cfg, err := awsconfig.LoadDefaultConfig(context.Background(), + awsconfig.WithRegion(region), + awsconfig.WithCredentialsProvider(credentials.NewStaticCredentialsProvider(accessKey, secretKey, "")), + awsconfig.WithEndpointResolverWithOptions(customResolver), + ) + if err != nil { + log.Error("failed to create AWS config", "err", err) + os.Exit(1) + } + + client := sqs.NewFromConfig(cfg) + ctx := context.Background() + urlResult, err := client.GetQueueUrl(ctx, &sqs.GetQueueUrlInput{QueueName: aws.String(queueName)}) + if err != nil { + log.Error("GetQueueUrl failed", "queue", queueName, "err", err) + os.Exit(1) + } + queueURL := aws.ToString(urlResult.QueueUrl) + log.Info("sqs-consumer started", "queue", queueName, "endpoint", endpoint, "functionURL", functionURL) + + httpClient := &http.Client{Timeout: 30 * time.Second} + tm := newTokenManager(httpClient, functionURL, routerLoginURL, routerUser, routerPass, staticToken) + + ticker := time.NewTicker(time.Duration(pollSec) * time.Second) + defer ticker.Stop() + + for range ticker.C { + poll(ctx, log, client, httpClient, tm, queueURL, functionURL, maxMsg, maxRetries) + } +} + +func poll(ctx context.Context, log *slog.Logger, sqsClient *sqs.Client, httpClient *http.Client, + tm *tokenManager, queueURL, functionURL string, maxMsg int32, maxRetries int) { + + result, err := sqsClient.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{ + QueueUrl: aws.String(queueURL), + MaxNumberOfMessages: maxMsg, + WaitTimeSeconds: 5, + }) + if err != nil { + log.Error("SQS ReceiveMessage failed", "err", err) + return + } + + for _, msg := range result.Messages { + body := aws.ToString(msg.Body) + receipt := aws.ToString(msg.ReceiptHandle) + if invokeWithRetry(log, httpClient, tm, functionURL, body, maxRetries) { + if _, err := sqsClient.DeleteMessage(ctx, &sqs.DeleteMessageInput{ + QueueUrl: aws.String(queueURL), + ReceiptHandle: aws.String(receipt), + }); err != nil { + log.Error("DeleteMessage failed", "err", err) + } else { + log.Info("message processed", "msgId", aws.ToString(msg.MessageId)) + } + } + } +} + +func invokeWithRetry(log *slog.Logger, client *http.Client, tm *tokenManager, url, body string, maxRetries int) bool { + for attempt := 1; attempt <= maxRetries; attempt++ { + token := tm.getToken(log) + req, _ := http.NewRequest(http.MethodPost, url, bytes.NewBufferString(body)) + req.Header.Set("Content-Type", "application/json") + if token != "" { + req.Header.Set("Authorization", "Bearer "+token) + } + resp, err := client.Do(req) + if err != nil { + log.Warn("function invoke error", "attempt", attempt, "err", err) + time.Sleep(time.Duration(attempt) * time.Second) + continue + } + io.Copy(io.Discard, resp.Body) + resp.Body.Close() + if resp.StatusCode >= 200 && resp.StatusCode < 300 { + return true + } + if resp.StatusCode == http.StatusUnauthorized { + log.Warn("got 401, refreshing token", "attempt", attempt) + tm.invalidate() + } else { + log.Warn("function returned non-2xx", "attempt", attempt, "status", resp.StatusCode) + } + time.Sleep(time.Duration(attempt) * time.Second) + } + log.Error("all retries exhausted, skipping message", "url", url) + return false +} + +func getenv(key, def string) string { + if v := os.Getenv(key); v != "" { + return v + } + return def +} + +func mustenv(key string) string { + v := os.Getenv(key) + if v == "" { + fmt.Fprintf(os.Stderr, "ERROR: env var %s is required\n", key) + os.Exit(1) + } + return v +} + +func parseInt(s string, def int) int { + n, err := strconv.Atoi(s) + if err != nil || n <= 0 { + return def + } + return n +} diff --git a/terraform/provider/internal/client/client.go b/terraform/provider/internal/client/client.go index 6da00b7..b20a0ea 100644 --- a/terraform/provider/internal/client/client.go +++ b/terraform/provider/internal/client/client.go @@ -302,8 +302,8 @@ func (c *Client) CreateMQTrigger(ctx context.Context, name, ns, functionName, qu deployName := mqDeployName(name) labels := map[string]string{ - mqManagedByLabel: mqManagedByVal, - mqComponentLabel: mqComponentVal, + mqManagedByLabel: mqManagedByVal, + mqComponentLabel: mqComponentVal, "mq-trigger-name": name, }