- api/v1alpha1/service_types.go: убрать +kubebuilder:default=30
- invoke.go: TimeoutSec=0 → &http.Client{} (без таймаута)
- services.go: валидация timeout_sec < 0 || > 900 → HTTP 400
- service_resource.go: TF schema Optional (без Computed); 0 → Int64Null()
- deployments/k8s/operator.yaml: v0.1.47 → v0.1.48
- doc/: progress.md + api/design.md (модель Service) + decisions/log.md
- examples/POSTGRES/: bug_hunter.sh, chaos_marathon.sh, chaos_marathon.tf
245 lines
9.7 KiB
Go
245 lines
9.7 KiB
Go
// Изменено: 2026-03-20 (function-service-split: dual-mode invoke)
|
|
// invoke.go — обработчик вызова функций и сервисов.
|
|
// Маршрут: ANY /fn/{namespace}/{name} и /fn/{namespace}/{name}/**
|
|
// Не защищён auth-токеном — публичный эндпоинт.
|
|
//
|
|
// Два режима:
|
|
// 1. Service mode (sless_service): проксирует запрос к ClusterIP Deployment-пода.
|
|
// URL: http://{name}.sless-fn-{namespace}.svc.cluster.local:8080
|
|
// Таймаут = Spec.TimeoutSec + 5s буфер.
|
|
//
|
|
// 2. Function mode (sless_function): создаёт FunctionJob CRD, ждёт завершения,
|
|
// возвращает содержимое Message (stdout функции).
|
|
// Таймаут = Spec.TimeoutSec (or 30s default).
|
|
//
|
|
// Порядок поиска: сначала Service CRD → Function CRD → 404.
|
|
|
|
package handler
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/gorilla/mux"
|
|
"k8s.io/apimachinery/pkg/api/errors"
|
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
|
|
|
slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1"
|
|
)
|
|
|
|
// 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 под конкретный вызов.
|
|
// timeoutSec == 0 → Timeout: 0 = без ограничения (пользователь не установил лимит).
|
|
// timeoutSec > 0 → Timeout: timeoutSec + 5s буфер на сетевые задержки.
|
|
func invokeHTTPClient(timeoutSec int32) *http.Client {
|
|
if timeoutSec <= 0 {
|
|
// http.Client{Timeout: 0} в Go означает отсутствие таймаута.
|
|
return &http.Client{}
|
|
}
|
|
return &http.Client{Timeout: time.Duration(timeoutSec)*time.Second + 5*time.Second}
|
|
}
|
|
|
|
// InvokeFunction обрабатывает вызов /fn/{namespace}/{name}[/**].
|
|
// Определяет режим по типу ресурса: Service (proxy) или Function (Job).
|
|
func (h *Handler) InvokeFunction(w http.ResponseWriter, r *http.Request) {
|
|
vars := mux.Vars(r)
|
|
ns := vars["namespace"]
|
|
name := vars["name"]
|
|
|
|
// Пробуем Service CRD первым — это основной режим long-running функций
|
|
svc := &slessv1alpha1.Service{}
|
|
if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, svc); err == nil {
|
|
h.invokeServiceProxy(w, r, ns, name, svc.Spec.TimeoutSec)
|
|
return
|
|
} else if !errors.IsNotFound(err) {
|
|
h.Log.Error("invoke: get service CRD", "err", err, "ns", ns, "name", name)
|
|
writeJSON(w, http.StatusInternalServerError, errResp("internal error"))
|
|
return
|
|
}
|
|
|
|
// Пробуем Function CRD — oneshot режим (create Job, wait, return result)
|
|
fn := &slessv1alpha1.Function{}
|
|
if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, fn); err == nil {
|
|
h.invokeFunctionJob(w, r, ns, name, fn)
|
|
return
|
|
} else if !errors.IsNotFound(err) {
|
|
h.Log.Error("invoke: get function CRD", "err", err, "ns", ns, "name", name)
|
|
writeJSON(w, http.StatusInternalServerError, errResp("internal error"))
|
|
return
|
|
}
|
|
|
|
writeJSON(w, http.StatusNotFound, errResp("function or service not found"))
|
|
}
|
|
|
|
// invokeServiceProxy проксирует запрос к ClusterIP сервиса в кластере.
|
|
// Service — long-running Deployment, постоянно доступен по внутреннему DNS.
|
|
func (h *Handler) invokeServiceProxy(w http.ResponseWriter, r *http.Request, ns, name string, timeoutSec int32) {
|
|
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 к k8s Service (DNS внутри кластера)
|
|
target := fmt.Sprintf("http://%s.sless-fn-%s.svc.cluster.local:8080%s", name, ns, subPath)
|
|
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, "name", name)
|
|
writeJSON(w, http.StatusInternalServerError, errResp("failed to create proxy request"))
|
|
return
|
|
}
|
|
|
|
// Content-Type и Content-Length обязательны для корректной работы Python/Node серверов.
|
|
// 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" — k8s Service не существует (сервис удалён или не задеплоен)
|
|
if strings.Contains(err.Error(), "no such host") {
|
|
writeJSON(w, http.StatusNotFound, errResp("service not found or not ready"))
|
|
return
|
|
}
|
|
h.Log.Error("invoke: service unreachable", "err", err, "ns", ns, "name", name, "target", target)
|
|
writeJSON(w, http.StatusBadGateway, errResp("service unreachable: "+err.Error()))
|
|
return
|
|
}
|
|
defer resp.Body.Close()
|
|
|
|
// Пробрасываем заголовки из ответа функции, фильтруя hop-by-hop
|
|
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)
|
|
}
|
|
|
|
// invokeFunctionJob создаёт FunctionJob CRD и синхронно ждёт завершения (polling 2s).
|
|
// Предназначен для sless_function — oneshot вызов без постоянного пода.
|
|
// Тело запроса передаётся как EventJSON в FunctionJobSpec.
|
|
// Job удаляется после получения результата (best-effort cleanup).
|
|
func (h *Handler) invokeFunctionJob(w http.ResponseWriter, r *http.Request, ns, name string, fn *slessv1alpha1.Function) {
|
|
if fn.Status.Phase != slessv1alpha1.FunctionPhaseReady {
|
|
writeJSON(w, http.StatusServiceUnavailable, errResp(
|
|
fmt.Sprintf("function not ready (phase: %s)", fn.Status.Phase),
|
|
))
|
|
return
|
|
}
|
|
|
|
// Читаем тело запроса как EventJSON (максимум 1MB).
|
|
// Функция получит это в handle(event) через runner.
|
|
eventJSON := "{}"
|
|
if r.Body != nil {
|
|
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20))
|
|
if err == nil && len(body) > 0 {
|
|
eventJSON = string(body)
|
|
}
|
|
}
|
|
|
|
// Уникальное имя Job = имя функции + unix nanos для уникальности
|
|
rawName := fmt.Sprintf("%s-%d", name, time.Now().UnixNano())
|
|
jobName := rawName
|
|
if len(jobName) > 63 {
|
|
jobName = jobName[:63]
|
|
}
|
|
|
|
// Копируем поля из Function CRD — FunctionJob теперь самодостаточен (FunctionRef удалён).
|
|
fj := &slessv1alpha1.FunctionJob{
|
|
ObjectMeta: metav1.ObjectMeta{
|
|
Name: jobName,
|
|
Namespace: ns,
|
|
},
|
|
Spec: slessv1alpha1.FunctionJobSpec{
|
|
Runtime: fn.Spec.Runtime,
|
|
Entrypoint: fn.Spec.Entrypoint,
|
|
S3Bucket: fn.Spec.S3Bucket,
|
|
S3Key: fn.Spec.S3Key,
|
|
MemoryMB: fn.Spec.MemoryMB,
|
|
TimeoutSec: fn.Spec.TimeoutSec,
|
|
Env: fn.Spec.Env,
|
|
EventJSON: eventJSON,
|
|
RunID: time.Now().UnixNano(),
|
|
},
|
|
}
|
|
if err := h.K8s.Create(r.Context(), fj); err != nil {
|
|
h.Log.Error("invoke: create FunctionJob", "err", err, "ns", ns, "fn", name)
|
|
writeJSON(w, http.StatusInternalServerError, errResp("failed to create job: "+err.Error()))
|
|
return
|
|
}
|
|
|
|
// Cleanup после получения результата — best-effort, не блокирует ответ.
|
|
// Используем context.Background() т.к. r.Context() может быть уже закрыт.
|
|
defer func() {
|
|
delCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
_ = h.K8s.Delete(delCtx, fj)
|
|
}()
|
|
|
|
// Опрашиваем каждые 2 секунды пока не завершится или не истечёт таймаут
|
|
timeoutSec := fn.Spec.TimeoutSec
|
|
if timeoutSec <= 0 {
|
|
timeoutSec = 30
|
|
}
|
|
deadline := time.Now().Add(time.Duration(timeoutSec) * time.Second)
|
|
|
|
for time.Now().Before(deadline) {
|
|
select {
|
|
case <-r.Context().Done():
|
|
writeJSON(w, http.StatusGatewayTimeout, errResp("request cancelled"))
|
|
return
|
|
case <-time.After(2 * time.Second):
|
|
}
|
|
|
|
current := &slessv1alpha1.FunctionJob{}
|
|
if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: jobName, Namespace: ns}, current); err != nil {
|
|
h.Log.Error("invoke: poll FunctionJob", "err", err, "job", jobName)
|
|
continue
|
|
}
|
|
|
|
switch current.Status.Phase {
|
|
case slessv1alpha1.FunctionJobPhaseSucceeded:
|
|
writeJSON(w, http.StatusOK, map[string]string{"result": current.Status.Message})
|
|
return
|
|
case slessv1alpha1.FunctionJobPhaseFailed:
|
|
writeJSON(w, http.StatusInternalServerError, errResp(current.Status.Message))
|
|
return
|
|
}
|
|
}
|
|
|
|
writeJSON(w, http.StatusGatewayTimeout, errResp(
|
|
fmt.Sprintf("function %s/%s timed out after %ds", ns, name, timeoutSec),
|
|
))
|
|
}
|