From 2ee9cae6d233cc47dfa157aeb21c745ee055fda3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E2=80=9CNaeel=E2=80=9D?= Date: Sat, 7 Mar 2026 18:36:03 +0400 Subject: [PATCH] =?UTF-8?q?feat:=20proxy=20/fn/{namespace}/{name}=20?= =?UTF-8?q?=E2=80=94=20=D0=BE=D0=B1=D1=85=D0=BE=D0=B4=20wildcard=20DNS?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Проблема: wildcard DNS *.fn.kube5s.ru недоступен. Решение: прокси через sless-api.kube5s.ru/fn/{ns}/{name}. - handler/invoke.go: прокси к Service функции внутри кластера - router.go: /fn/ без auth токена, /v1/ с auth (gorilla Use()) - config.go: поле ExternalURL (EXTERNAL_URL env) - trigger_controller.go: если ExternalURL задан — URL = ExternalURL/fn/{ns}/{fn} иначе fallback: Ingress + поддомен (прежнее поведение) - operator.yaml: EXTERNAL_URL=https://sless-api.kube5s.ru, image v0.1.5 Оператор v0.1.5 задеплоен. E2E: curl https://sless-api.kube5s.ru/fn/default/hello-node → {"message":"Hello, Naeel! (nodejs20)"} --- controllers/trigger_controller.go | 429 +++++++++++++++--------------- deployments/k8s/operator.yaml | 5 +- internal/api/handler/invoke.go | 65 +++++ internal/api/router.go | 19 +- internal/config/config.go | 14 +- main.go | 1 + 6 files changed, 313 insertions(+), 220 deletions(-) create mode 100644 internal/api/handler/invoke.go diff --git a/controllers/trigger_controller.go b/controllers/trigger_controller.go index 9c002f0..5888054 100644 --- a/controllers/trigger_controller.go +++ b/controllers/trigger_controller.go @@ -6,31 +6,32 @@ package controllers import ( -"context" -"fmt" + "context" + "fmt" -batchv1 "k8s.io/api/batch/v1" -corev1 "k8s.io/api/core/v1" -netv1 "k8s.io/api/networking/v1" -"k8s.io/apimachinery/pkg/api/errors" -metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" -"k8s.io/apimachinery/pkg/runtime" -ctrl "sigs.k8s.io/controller-runtime" -"sigs.k8s.io/controller-runtime/pkg/client" -"sigs.k8s.io/controller-runtime/pkg/log" + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + netv1 "k8s.io/api/networking/v1" + "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/log" -slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1" + slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1" ) // TriggerReconciler reconciles a Trigger object type TriggerReconciler struct { -client.Client -Scheme *runtime.Scheme -IngressHost string // базовый домен для http триггеров, например: fn.kube5s.ru + client.Client + Scheme *runtime.Scheme + IngressHost string // базовый домен для http триггеров (fallback): fn.kube5s.ru + // ExternalURL — если задан, URL функции = ExternalURL/fn/{namespace}/{name} + // (обходит wildcard DNS, трафик идёт через sless-api.kube5s.ru) + ExternalURL string } -//+kubebuilder:rbac:groups=sless.kube5s.ru,resources=triggers,verbs=get;list;watch;create;update;patch;delete -//+kubebuilder:rbac:groups=sless.kube5s.ru,resources=triggers/status,verbs=get;update;patch //+kubebuilder:rbac:groups=sless.kube5s.ru,resources=triggers/finalizers,verbs=update //+kubebuilder:rbac:groups=networking.k8s.io,resources=ingresses,verbs=get;list;watch;create;update;patch;delete //+kubebuilder:rbac:groups="",resources=services,verbs=get;list;watch;create;update;patch;delete @@ -40,221 +41,229 @@ const triggerFinalizer = "sless.kube5s.ru/trigger-finalizer" // Reconcile — управляет ресурсами для триггера в зависимости от его типа. func (r *TriggerReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) { -logger := log.FromContext(ctx) + logger := log.FromContext(ctx) -tr := &slessv1alpha1.Trigger{} -if err := r.Get(ctx, req.NamespacedName, tr); err != nil { -if errors.IsNotFound(err) { -return ctrl.Result{}, nil -} -return ctrl.Result{}, fmt.Errorf("get trigger: %w", err) + tr := &slessv1alpha1.Trigger{} + if err := r.Get(ctx, req.NamespacedName, tr); err != nil { + if errors.IsNotFound(err) { + return ctrl.Result{}, nil + } + return ctrl.Result{}, fmt.Errorf("get trigger: %w", err) + } + + // Обработка удаления + if !tr.DeletionTimestamp.IsZero() { + return r.handleTriggerDeletion(ctx, tr) + } + + // Добавляем finalizer если ещё нет + if !containsString(tr.Finalizers, triggerFinalizer) { + tr.Finalizers = append(tr.Finalizers, triggerFinalizer) + if err := r.Update(ctx, tr); err != nil { + return ctrl.Result{}, fmt.Errorf("add trigger finalizer: %w", err) + } + return ctrl.Result{Requeue: true}, nil + } + + // Проверяем что Function существует и готова + fn := &slessv1alpha1.Function{} + if err := r.Get(ctx, client.ObjectKey{Name: tr.Spec.FunctionRef, Namespace: tr.Namespace}, fn); err != nil { + if errors.IsNotFound(err) { + tr.Status.Active = false + tr.Status.Message = "function not found: " + tr.Spec.FunctionRef + _ = r.Status().Update(ctx, tr) + return ctrl.Result{}, nil + } + return ctrl.Result{}, fmt.Errorf("get function: %w", err) + } + if fn.Status.Phase != slessv1alpha1.FunctionPhaseReady { + tr.Status.Active = false + tr.Status.Message = "waiting for function to be Ready (current: " + string(fn.Status.Phase) + ")" + _ = r.Status().Update(ctx, tr) + // Повторный reconcile придёт когда Function изменится (watch ниже) + return ctrl.Result{}, nil + } + + switch tr.Spec.Type { + case slessv1alpha1.TriggerTypeHTTP: + logger.Info("reconcile http trigger", "trigger", tr.Name) + return r.reconcileHTTP(ctx, tr, fn) + case slessv1alpha1.TriggerTypeCron: + logger.Info("reconcile cron trigger", "trigger", tr.Name) + return r.reconcileCron(ctx, tr, fn) + } + + return ctrl.Result{}, nil } -// Обработка удаления -if !tr.DeletionTimestamp.IsZero() { -return r.handleTriggerDeletion(ctx, tr) -} - -// Добавляем finalizer если ещё нет -if !containsString(tr.Finalizers, triggerFinalizer) { -tr.Finalizers = append(tr.Finalizers, triggerFinalizer) -if err := r.Update(ctx, tr); err != nil { -return ctrl.Result{}, fmt.Errorf("add trigger finalizer: %w", err) -} -return ctrl.Result{Requeue: true}, nil -} - -// Проверяем что Function существует и готова -fn := &slessv1alpha1.Function{} -if err := r.Get(ctx, client.ObjectKey{Name: tr.Spec.FunctionRef, Namespace: tr.Namespace}, fn); err != nil { -if errors.IsNotFound(err) { -tr.Status.Active = false -tr.Status.Message = "function not found: " + tr.Spec.FunctionRef -_ = r.Status().Update(ctx, tr) -return ctrl.Result{}, nil -} -return ctrl.Result{}, fmt.Errorf("get function: %w", err) -} -if fn.Status.Phase != slessv1alpha1.FunctionPhaseReady { -tr.Status.Active = false -tr.Status.Message = "waiting for function to be Ready (current: " + string(fn.Status.Phase) + ")" -_ = r.Status().Update(ctx, tr) -// Повторный reconcile придёт когда Function изменится (watch ниже) -return ctrl.Result{}, nil -} - -switch tr.Spec.Type { -case slessv1alpha1.TriggerTypeHTTP: -logger.Info("reconcile http trigger", "trigger", tr.Name) -return r.reconcileHTTP(ctx, tr, fn) -case slessv1alpha1.TriggerTypeCron: -logger.Info("reconcile cron trigger", "trigger", tr.Name) -return r.reconcileCron(ctx, tr, fn) -} - -return ctrl.Result{}, nil -} - -// reconcileHTTP создаёт Service и Ingress для HTTP триггера. -// URL функции: https://{funcName}-{namespace}.{IngressHost} +// reconcileHTTP создаёт Service для HTTP триггера. +// Если ExternalURL задан — URL = {ExternalURL}/fn/{namespace}/{name} (прокси через API). +// Иначе (fallback) — URL = https://{funcName}-{namespace}.{IngressHost} + создаёт Ingress. func (r *TriggerReconciler) reconcileHTTP(ctx context.Context, tr *slessv1alpha1.Trigger, fn *slessv1alpha1.Function) (ctrl.Result, error) { -deployNS := "sless-fn-" + tr.Namespace -host := fmt.Sprintf("%s-%s.%s", fn.Name, tr.Namespace, r.IngressHost) -pathType := netv1.PathTypePrefix + deployNS := "sless-fn-" + tr.Namespace -// Service — направляет трафик к Deployment функции -wantSvc := &corev1.Service{ -ObjectMeta: metav1.ObjectMeta{Name: fn.Name, Namespace: deployNS}, -Spec: corev1.ServiceSpec{ -Selector: map[string]string{"app": fn.Name}, -Ports: []corev1.ServicePort{{Port: 8080, Protocol: corev1.ProtocolTCP}}, -}, -} -existingSvc := &corev1.Service{} -if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, existingSvc); err != nil { -if errors.IsNotFound(err) { -if err := r.Create(ctx, wantSvc); err != nil { -return ctrl.Result{}, fmt.Errorf("create service: %w", err) -} -} else { -return ctrl.Result{}, fmt.Errorf("get service: %w", err) -} -} + // Service — направляет трафик к Deployment функции (нужен и для прокси, и для Ingress) + wantSvc := &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{Name: fn.Name, Namespace: deployNS}, + Spec: corev1.ServiceSpec{ + Selector: map[string]string{"app": fn.Name}, + Ports: []corev1.ServicePort{{Port: 8080, Protocol: corev1.ProtocolTCP}}, + }, + } + existingSvc := &corev1.Service{} + if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, existingSvc); err != nil { + if errors.IsNotFound(err) { + if err := r.Create(ctx, wantSvc); err != nil { + return ctrl.Result{}, fmt.Errorf("create service: %w", err) + } + } else { + return ctrl.Result{}, fmt.Errorf("get service: %w", err) + } + } -// Ingress — внешний доступ по host -wantIng := &netv1.Ingress{ -ObjectMeta: metav1.ObjectMeta{ -Name: fn.Name, -Namespace: deployNS, -Annotations: map[string]string{ -"kubernetes.io/ingress.class": "nginx", -}, -}, -Spec: netv1.IngressSpec{ -Rules: []netv1.IngressRule{ -{ -Host: host, -IngressRuleValue: netv1.IngressRuleValue{ -HTTP: &netv1.HTTPIngressRuleValue{ -Paths: []netv1.HTTPIngressPath{ -{ -Path: "/", -PathType: &pathType, -Backend: netv1.IngressBackend{ -Service: &netv1.IngressServiceBackend{ -Name: fn.Name, -Port: netv1.ServiceBackendPort{Number: 8080}, -}, -}, -}, -}, -}, -}, -}, -}, -}, -} -existingIng := &netv1.Ingress{} -if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, existingIng); err != nil { -if errors.IsNotFound(err) { -if err := r.Create(ctx, wantIng); err != nil { -return ctrl.Result{}, fmt.Errorf("create ingress: %w", err) -} -} else { -return ctrl.Result{}, fmt.Errorf("get ingress: %w", err) -} -} + var funcURL string + if r.ExternalURL != "" { + // Режим прокси через API сервер — не нужен wildcard DNS + funcURL = fmt.Sprintf("%s/fn/%s/%s", r.ExternalURL, tr.Namespace, fn.Name) + } else { + // Fallback: Ingress с поддоменом (требует wildcard DNS *.IngressHost) + host := fmt.Sprintf("%s-%s.%s", fn.Name, tr.Namespace, r.IngressHost) + pathType := netv1.PathTypePrefix + wantIng := &netv1.Ingress{ + ObjectMeta: metav1.ObjectMeta{ + Name: fn.Name, + Namespace: deployNS, + Annotations: map[string]string{ + "kubernetes.io/ingress.class": "nginx", + }, + }, + Spec: netv1.IngressSpec{ + Rules: []netv1.IngressRule{ + { + Host: host, + IngressRuleValue: netv1.IngressRuleValue{ + HTTP: &netv1.HTTPIngressRuleValue{ + Paths: []netv1.HTTPIngressPath{ + { + Path: "/", + PathType: &pathType, + Backend: netv1.IngressBackend{ + Service: &netv1.IngressServiceBackend{ + Name: fn.Name, + Port: netv1.ServiceBackendPort{Number: 8080}, + }, + }, + }, + }, + }, + }, + }, + }, + }, + } + existingIng := &netv1.Ingress{} + if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, existingIng); err != nil { + if errors.IsNotFound(err) { + if err := r.Create(ctx, wantIng); err != nil { + return ctrl.Result{}, fmt.Errorf("create ingress: %w", err) + } + } else { + return ctrl.Result{}, fmt.Errorf("get ingress: %w", err) + } + } + funcURL = "https://" + host + } -tr.Status.Active = true -tr.Status.URL = "https://" + host -tr.Status.Message = "" -if err := r.Status().Update(ctx, tr); err != nil { -return ctrl.Result{}, fmt.Errorf("update trigger status: %w", err) -} -return ctrl.Result{}, nil + tr.Status.Active = true + tr.Status.URL = funcURL + tr.Status.Message = "" + if err := r.Status().Update(ctx, tr); err != nil { + return ctrl.Result{}, fmt.Errorf("update trigger status: %w", err) + } + return ctrl.Result{}, nil } // reconcileCron создаёт CronJob который вызывает функцию по HTTP внутри кластера. // curl делает POST на внутренний Service функции — это исключает внешний round-trip. func (r *TriggerReconciler) reconcileCron(ctx context.Context, tr *slessv1alpha1.Trigger, fn *slessv1alpha1.Function) (ctrl.Result, error) { -deployNS := "sless-fn-" + tr.Namespace -// Внутренний URL: Service должен быть создан HTTP триггером или заранее -funcURL := fmt.Sprintf("http://%s.%s.svc.cluster.local:8080", fn.Name, deployNS) + deployNS := "sless-fn-" + tr.Namespace + // Внутренний URL: Service должен быть создан HTTP триггером или заранее + funcURL := fmt.Sprintf("http://%s.%s.svc.cluster.local:8080", fn.Name, deployNS) -wantCJ := &batchv1.CronJob{ -ObjectMeta: metav1.ObjectMeta{ -Name: tr.Name, -Namespace: tr.Namespace, -Labels: map[string]string{"managed-by": "sless", "trigger": tr.Name}, -}, -Spec: batchv1.CronJobSpec{ -Schedule: tr.Spec.Schedule, -JobTemplate: batchv1.JobTemplateSpec{ -Spec: batchv1.JobSpec{ -Template: corev1.PodTemplateSpec{ -Spec: corev1.PodSpec{ -RestartPolicy: corev1.RestartPolicyOnFailure, -Containers: []corev1.Container{ -{ -// curlimages/curl вызывает функцию по внутреннему адресу -Name: "invoker", -Image: "curlimages/curl:latest", -Command: []string{"curl", "-sf", "-X", "POST", funcURL}, -}, -}, -}, -}, -}, -}, -}, -} + wantCJ := &batchv1.CronJob{ + ObjectMeta: metav1.ObjectMeta{ + Name: tr.Name, + Namespace: tr.Namespace, + Labels: map[string]string{"managed-by": "sless", "trigger": tr.Name}, + }, + Spec: batchv1.CronJobSpec{ + Schedule: tr.Spec.Schedule, + JobTemplate: batchv1.JobTemplateSpec{ + Spec: batchv1.JobSpec{ + Template: corev1.PodTemplateSpec{ + Spec: corev1.PodSpec{ + RestartPolicy: corev1.RestartPolicyOnFailure, + Containers: []corev1.Container{ + { + // curlimages/curl вызывает функцию по внутреннему адресу + Name: "invoker", + Image: "curlimages/curl:latest", + Command: []string{"curl", "-sf", "-X", "POST", funcURL}, + }, + }, + }, + }, + }, + }, + }, + } -existing := &batchv1.CronJob{} -if err := r.Get(ctx, client.ObjectKey{Name: tr.Name, Namespace: tr.Namespace}, existing); err != nil { -if errors.IsNotFound(err) { -if err := r.Create(ctx, wantCJ); err != nil { -return ctrl.Result{}, fmt.Errorf("create cronjob: %w", err) -} -} else { -return ctrl.Result{}, fmt.Errorf("get cronjob: %w", err) -} -} else { -// Синхронизируем расписание если оно изменилось -existing.Spec.Schedule = tr.Spec.Schedule -if err := r.Update(ctx, existing); err != nil { -return ctrl.Result{}, fmt.Errorf("update cronjob schedule: %w", err) -} -} + existing := &batchv1.CronJob{} + if err := r.Get(ctx, client.ObjectKey{Name: tr.Name, Namespace: tr.Namespace}, existing); err != nil { + if errors.IsNotFound(err) { + if err := r.Create(ctx, wantCJ); err != nil { + return ctrl.Result{}, fmt.Errorf("create cronjob: %w", err) + } + } else { + return ctrl.Result{}, fmt.Errorf("get cronjob: %w", err) + } + } else { + // Синхронизируем расписание если оно изменилось + existing.Spec.Schedule = tr.Spec.Schedule + if err := r.Update(ctx, existing); err != nil { + return ctrl.Result{}, fmt.Errorf("update cronjob schedule: %w", err) + } + } -tr.Status.Active = true -tr.Status.Message = "" -if err := r.Status().Update(ctx, tr); err != nil { -return ctrl.Result{}, fmt.Errorf("update trigger status: %w", err) -} -return ctrl.Result{}, nil + tr.Status.Active = true + tr.Status.Message = "" + if err := r.Status().Update(ctx, tr); err != nil { + return ctrl.Result{}, fmt.Errorf("update trigger status: %w", err) + } + return ctrl.Result{}, nil } // handleTriggerDeletion удаляет ресурсы триггера и убирает finalizer. func (r *TriggerReconciler) handleTriggerDeletion(ctx context.Context, tr *slessv1alpha1.Trigger) (ctrl.Result, error) { -if tr.Spec.Type == slessv1alpha1.TriggerTypeCron { -cj := &batchv1.CronJob{} -if err := r.Get(ctx, client.ObjectKey{Name: tr.Name, Namespace: tr.Namespace}, cj); err == nil { -_ = r.Delete(ctx, cj) -} -} -// Ingress/Service для HTTP триггеров оставляем — они могут быть нужны другим триггерам, -// полная очистка происходит при удалении Function через function_controller finalizer -tr.Finalizers = removeString(tr.Finalizers, triggerFinalizer) -if err := r.Update(ctx, tr); err != nil { -return ctrl.Result{}, fmt.Errorf("remove trigger finalizer: %w", err) -} -return ctrl.Result{}, nil + if tr.Spec.Type == slessv1alpha1.TriggerTypeCron { + cj := &batchv1.CronJob{} + if err := r.Get(ctx, client.ObjectKey{Name: tr.Name, Namespace: tr.Namespace}, cj); err == nil { + _ = r.Delete(ctx, cj) + } + } + // Ingress/Service для HTTP триггеров оставляем — они могут быть нужны другим триггерам, + // полная очистка происходит при удалении Function через function_controller finalizer + tr.Finalizers = removeString(tr.Finalizers, triggerFinalizer) + if err := r.Update(ctx, tr); err != nil { + return ctrl.Result{}, fmt.Errorf("remove trigger finalizer: %w", err) + } + return ctrl.Result{}, nil } // SetupWithManager регистрирует контроллер и настраивает watch на Function. // Когда Function переходит в Ready — Trigger автоматически пересчитывается. func (r *TriggerReconciler) SetupWithManager(mgr ctrl.Manager) error { -return ctrl.NewControllerManagedBy(mgr). -For(&slessv1alpha1.Trigger{}). -Complete(r) + return ctrl.NewControllerManagedBy(mgr). + For(&slessv1alpha1.Trigger{}). + Complete(r) } diff --git a/deployments/k8s/operator.yaml b/deployments/k8s/operator.yaml index 3112d3b..e515874 100644 --- a/deployments/k8s/operator.yaml +++ b/deployments/k8s/operator.yaml @@ -31,6 +31,9 @@ data: REGISTRY_SECRET: "sless-registry-auth" API_PORT: "9090" INGRESS_HOST: "fn.kube5s.ru" + # EXTERNAL_URL — если задан, URL функции = EXTERNAL_URL/fn/{namespace}/{name} + # Позволяет обойтись без wildcard DNS *.fn.kube5s.ru + EXTERNAL_URL: "https://sless-api.kube5s.ru" --- # Secret создаётся отдельно через kubectl (не коммитить секреты в git!) # Описание ключей: @@ -67,7 +70,7 @@ spec: containers: - name: operator # При обновлении версии оператора — менять тег здесь (не latest!) - image: naeel/sless-operator:v0.1.4 + image: naeel/sless-operator:v0.1.5 # Always — чтобы всегда тянуть по точному тегу (не кешировать старый) imagePullPolicy: Always ports: diff --git a/internal/api/handler/invoke.go b/internal/api/handler/invoke.go new file mode 100644 index 0000000..aa86588 --- /dev/null +++ b/internal/api/handler/invoke.go @@ -0,0 +1,65 @@ +// Изменено: 2026-03-07 +// invoke.go — прокси-обработчик для вызова HTTP-триггеров функций. +// Маршрут: ANY /fn/{namespace}/{name} +// Не защищён auth-токеном — это публичный эндпоинт для вызова функций. +// Проксирует запрос к ClusterIP Service функции внутри кластера: +// http://{name}.sless-fn-{namespace}.svc.cluster.local:8080 +// Функция должна быть задеплоена (HTTP триггер активен). + +package handler + +import ( + "fmt" + "io" + "net/http" + "time" + + "github.com/gorilla/mux" +) + +// httpClient используется для обращения к функциям внутри кластера. +// Таймаут 30s — достаточно для холодного старта функции. +var httpClient = &http.Client{Timeout: 30 * time.Second} + +// InvokeFunction проксирует входящий запрос к Service функции в кластере. +// Namespace выбирается из пути, имя функции — тоже из пути. +// Сохраняет метод, тело, и Content-Type заголовок оригинального запроса. +func (h *Handler) InvokeFunction(w http.ResponseWriter, r *http.Request) { + vars := mux.Vars(r) + ns := vars["namespace"] + name := vars["name"] + + // Внутренний URL к Service функции (DNS внутри кластера) + // Namespace функций: sless-fn-{user-namespace} + target := fmt.Sprintf("http://%s.sless-fn-%s.svc.cluster.local:8080", name, ns) + + // Создаём проксируемый запрос с тем же методом и телом + proxyReq, err := http.NewRequestWithContext(r.Context(), r.Method, target, r.Body) + if err != nil { + h.Log.Error("invoke: failed to create proxy request", "err", err, "ns", ns, "fn", name) + writeJSON(w, http.StatusInternalServerError, errResp("failed to create proxy request")) + return + } + + // Пробрасываем Content-Type если есть + if ct := r.Header.Get("Content-Type"); ct != "" { + proxyReq.Header.Set("Content-Type", ct) + } + + resp, err := httpClient.Do(proxyReq) + if err != nil { + h.Log.Error("invoke: function unreachable", "err", err, "ns", ns, "fn", name, "target", target) + writeJSON(w, http.StatusBadGateway, errResp("function unreachable: "+err.Error())) + return + } + defer resp.Body.Close() + + // Копируем заголовки и статус из ответа функции + for k, vals := range resp.Header { + for _, v := range vals { + w.Header().Add(k, v) + } + } + w.WriteHeader(resp.StatusCode) + _, _ = io.Copy(w, resp.Body) +} diff --git a/internal/api/router.go b/internal/api/router.go index 70ebfa2..1681a75 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -17,9 +17,14 @@ import ( // NewRouter собирает gorilla/mux роутер со всеми маршрутами. // apiToken — статический Bearer-токен для v1 аутентификации. +// /fn/{namespace}/{name} — публичный прокси для вызова функций, без auth. func NewRouter(h *handler.Handler, apiToken string, log *slog.Logger) http.Handler { r := mux.NewRouter() + // Публичный прокси для вызова HTTP-триггеров — без auth токена + // Все HTTP методы разрешены (GET/POST/PUT/... — решает сама функция) + r.PathPrefix("/fn/{namespace}/{name}").HandlerFunc(h.InvokeFunction) + // Суброутер для /v1 — все маршруты API v1 := r.PathPrefix("/v1").Subrouter() @@ -47,10 +52,12 @@ func NewRouter(h *handler.Handler, apiToken string, log *slog.Logger) http.Handl v1.HandleFunc("/namespaces/{namespace}/jobs/{name}", h.GetJob).Methods(http.MethodGet) v1.HandleFunc("/namespaces/{namespace}/jobs/{name}", h.DeleteJob).Methods(http.MethodDelete) - // Цепочка middleware: logging → auth → router - // Порядок важен: сначала логируем (чтобы видеть все запросы включая отклонённые), - // затем проверяем авторизацию. - return middleware.Logging(log, - middleware.Auth(apiToken, log, r), - ) + // Цепочка middleware: logging → (auth только для /v1/) → router + // /fn/ — без auth, /v1/ — с auth. + // Используем gorilla/mux Use() чтобы auth применялся только к v1 суброутеру. + v1.Use(func(next http.Handler) http.Handler { + return middleware.Auth(apiToken, log, next) + }) + + return middleware.Logging(log, r) } diff --git a/internal/config/config.go b/internal/config/config.go index a3e634b..d551d66 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -41,10 +41,15 @@ type Config struct { // BuilderImage — образ для сборки функций (kaniko или buildah) BuilderImage string - // IngressHost — базовый домен для HTTP триггеров, например: fn.kube5s.ru - // URL функции будет: https://{funcName}-{namespace}.{IngressHost} + // IngressHost — базовый домен для HTTP триггеров (используется только + // если ExternalURL не задан, для обратной совместимости). IngressHost string + // ExternalURL — публичный URL API сервера, например: https://sless-api.kube5s.ru + // Если задан, URL функции формируется как: {ExternalURL}/fn/{namespace}/{name} + // Это позволяет обойтись без wildcard DNS *.fn. + ExternalURL string + // APIToken — статический Bearer-токен для v1 REST API аутентификации. // В prod заменить на вызов auth-сервиса. APIToken string @@ -117,12 +122,15 @@ func Load() (*Config, error) { cfg.BuilderImage = v } - // IngressHost — базовый домен для HTTP триггеров + // IngressHost — базовый домен для HTTP триггеров (fallback) cfg.IngressHost = os.Getenv("INGRESS_HOST") if cfg.IngressHost == "" { cfg.IngressHost = "fn.kube5s.ru" } + // ExternalURL — публичный URL сервиса для формирования URL функций + cfg.ExternalURL = os.Getenv("EXTERNAL_URL") + // APIToken — обязательный токен для REST API cfg.APIToken = os.Getenv("SLESS_API_TOKEN") if cfg.APIToken == "" { diff --git a/main.go b/main.go index d4e0583..286a268 100644 --- a/main.go +++ b/main.go @@ -140,6 +140,7 @@ func main() { Client: mgr.GetClient(), Scheme: mgr.GetScheme(), IngressHost: cfg.IngressHost, + ExternalURL: cfg.ExternalURL, }).SetupWithManager(mgr); err != nil { log.Error("unable to create controller", "controller", "Trigger", "err", err) os.Exit(1)