Files
sless/internal/api/handler/invoke.go
T
Naeel d7fda15d35 feat: Go runtime v0.1.1 (pgx/v5), stress-go-pgstorm, fix invoke.go dynamic timeout, nginx ingress timeout 900s
- runtimes/go1.23: добавлен pgx/v5 v5.7.2 в go.mod, сгенерирован go.sum
- runtimes/go1.23/Dockerfile: один stage golang:1.23-alpine, go mod download кеширует зависимости
- internal/builder/context.go: тег Go runtime v0.1.0 → v0.1.1
- internal/api/handler/invoke.go: таймаут прокси-клиента теперь динамический из Function.Spec.TimeoutSec + 5s буфер (был хардкод 30s)
- examples/POSTGRES/code/stress-go-pgstorm/handler.go: новая функция, 100 горутин, pgxpool, INSERT/COUNT/MAX, параметры: workers/duration_sec/max_delay_ms
- examples/POSTGRES/resources.tf: добавлены sless_function.stress_go_pgstorm + trigger, timeout_sec=700
- deployments/k8s/operator.yaml: nginx ingress proxy-read-timeout=900s, proxy-send-timeout=900s
- examples/*/main.tf: исправлен URL deck-api-test.ngcloud.ru → deck-test.ngcloud.ru (все 9 файлов)
- Оператор: v0.1.39 (pgx/v5) → v0.1.40 (dynamic timeout)
2026-03-19 21:33:41 +03:00

126 lines
5.6 KiB
Go

// Изменено: 2026-03-12
// invoke.go — прокси-обработчик для вызова HTTP-триггеров функций.
// Маршрут: ANY /fn/{namespace}/{name} и /fn/{namespace}/{name}/**
// Не защищён auth-токеном — это публичный эндпоинт для вызова функций.
// Проксирует запрос к ClusterIP Service функции внутри кластера:
// http://{name}.sless-fn-{namespace}.svc.cluster.local:8080
// Sub-path и query string пробрасываются как есть:
// /fn/ns/notes/add?title=x → http://notes.sless-fn-ns.svc.../add?title=x
// Таймаут берётся из Spec.TimeoutSec функции (+ 5s буфер) чтобы не резать
// длительные вызовы (stress-тесты, batch-задачи).
package handler
import (
"fmt"
"io"
"net/http"
"strings"
"time"
"github.com/gorilla/mux"
slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1"
"sigs.k8s.io/controller-runtime/pkg/client"
)
// hopByHopHeaders — заголовки которые нельзя пробрасывать через прокси (RFC 2616 §13.5.1).
// Они управляют соединением между двумя узлами, а не end-to-end.
// Особо опасен Transfer-Encoding: если пробросить его, клиент неверно интерпретирует тело.
var hopByHopHeaders = map[string]bool{
"Connection": true,
"Keep-Alive": true,
"Proxy-Authenticate": true,
"Proxy-Authorization": true,
"Te": true,
"Trailers": true,
"Transfer-Encoding": true,
"Upgrade": true,
}
// invokeHTTPClient создаёт http.Client с таймаутом под конкретный вызов.
// timeout = TimeoutSec функции + 5s буфер на сетевые задержки.
// Если TimeoutSec == 0 (не задан), используем 30s по умолчанию.
func invokeHTTPClient(timeoutSec int32) *http.Client {
t := time.Duration(timeoutSec)*time.Second + 5*time.Second
if timeoutSec <= 0 {
t = 30 * time.Second
}
return &http.Client{Timeout: t}
}
// InvokeFunction проксирует входящий запрос к Service функции в кластере.
// Namespace выбирается из пути, имя функции — тоже из пути.
// Сохраняет метод, тело, Content-Type, sub-path и query string.
// Таймаут прокси-клиента = Spec.TimeoutSec функции + 5s (резинка).
func (h *Handler) InvokeFunction(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
ns := vars["namespace"]
name := vars["name"]
// Смотрим TimeoutSec из Function CRD, чтобы не резать длительные вызовы.
// Если функция не найдена — продолжаем с дефолтным таймаутом (30s).
var timeoutSec int32
fn := &slessv1alpha1.Function{}
if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, fn); err == nil {
timeoutSec = fn.Spec.TimeoutSec
}
httpClient := invokeHTTPClient(timeoutSec)
// Вычисляем sub-path после /fn/{namespace}/{name}
// Например: /fn/default/notes/add → subPath = /add
prefix := "/fn/" + ns + "/" + name
subPath := strings.TrimPrefix(r.URL.Path, prefix)
// Внутренний URL к Service функции (DNS внутри кластера)
target := fmt.Sprintf("http://%s.sless-fn-%s.svc.cluster.local:8080%s", name, ns, subPath)
// Пробрасываем query string если есть
if r.URL.RawQuery != "" {
target += "?" + r.URL.RawQuery
}
// Создаём проксируемый запрос с тем же методом и телом
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 и Content-Length если есть.
// Content-Length обязателен: Python BaseHTTPRequestHandler читает тело
// ровно столько байт, сколько указано в заголовке; без него body = пусто.
if ct := r.Header.Get("Content-Type"); ct != "" {
proxyReq.Header.Set("Content-Type", ct)
}
proxyReq.ContentLength = r.ContentLength
resp, err := httpClient.Do(proxyReq)
if err != nil {
// "no such host" — Service не существует (функция удалена), возвращаем 404.
// Это отличается от временной сетевой ошибки: NXDOMAIN строго означает отсутствие записи.
if strings.Contains(err.Error(), "no such host") {
writeJSON(w, http.StatusNotFound, errResp("function not found"))
return
}
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()
// Копируем заголовки и статус из ответа функции.
// Hop-by-hop заголовки фильтруем: они управляют конкретным TCP-соединением
// и не должны пробрасываться через прокси (RFC 2616 §13.5.1).
for k, vals := range resp.Header {
if hopByHopHeaders[k] {
continue
}
for _, v := range vals {
w.Header().Add(k, v)
}
}
w.WriteHeader(resp.StatusCode)
_, _ = io.Copy(w, resp.Body)
}