287 lines
9.8 KiB
Go
287 lines
9.8 KiB
Go
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-<name>) с SQS credentials
|
||
// 2. Создаём K8s Deployment (mq-<name>) с образом naeel/sqs-consumer:v1.0
|
||
// FUNCTION_URL = http://router.fission.svc.cluster.local/<functionName>
|
||
// Лейблы: app.kubernetes.io/managed-by=fission-console, component=mq-trigger
|
||
//
|
||
// При LIST: deployments -n <ns> -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,
|
||
}
|
||
}
|