feat: sless_service CRD + ServiceReconciler, RBAC fix, split postgres/functions.tf, operator v0.1.41
This commit is contained in:
+208
-96
@@ -1,125 +1,237 @@
|
||||
// Изменено: 2026-03-12
|
||||
// invoke.go — прокси-обработчик для вызова HTTP-триггеров функций.
|
||||
// Изменено: 2026-03-20 (function-service-split: dual-mode invoke)
|
||||
// invoke.go — обработчик вызова функций и сервисов.
|
||||
// Маршрут: 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-задачи).
|
||||
// Не защищён 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 (
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
"context"
|
||||
"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"
|
||||
"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,
|
||||
"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 по умолчанию.
|
||||
// timeout = TimeoutSec + 5s буфер на сетевые задержки.
|
||||
// Если TimeoutSec == 0 (не задан), используем 35s по умолчанию.
|
||||
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}
|
||||
t := time.Duration(timeoutSec)*time.Second + 5*time.Second
|
||||
if timeoutSec <= 0 {
|
||||
t = 35 * time.Second
|
||||
}
|
||||
return &http.Client{Timeout: t}
|
||||
}
|
||||
|
||||
// InvokeFunction проксирует входящий запрос к Service функции в кластере.
|
||||
// Namespace выбирается из пути, имя функции — тоже из пути.
|
||||
// Сохраняет метод, тело, Content-Type, sub-path и query string.
|
||||
// Таймаут прокси-клиента = Spec.TimeoutSec функции + 5s (резинка).
|
||||
// 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"]
|
||||
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)
|
||||
// Пробуем 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
|
||||
}
|
||||
|
||||
// Вычисляем sub-path после /fn/{namespace}/{name}
|
||||
// Например: /fn/default/notes/add → subPath = /add
|
||||
prefix := "/fn/" + ns + "/" + name
|
||||
subPath := strings.TrimPrefix(r.URL.Path, prefix)
|
||||
// Пробуем 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
|
||||
}
|
||||
|
||||
// Внутренний 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
|
||||
writeJSON(w, http.StatusNotFound, errResp("function or service not found"))
|
||||
}
|
||||
|
||||
// Создаём проксируемый запрос с тем же методом и телом
|
||||
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
|
||||
// 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)
|
||||
}
|
||||
|
||||
// Пробрасываем 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
|
||||
// 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
|
||||
}
|
||||
|
||||
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()
|
||||
// Читаем тело запроса как 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)
|
||||
}
|
||||
}
|
||||
|
||||
// Копируем заголовки и статус из ответа функции.
|
||||
// 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)
|
||||
// Уникальное имя Job = имя функции + unix nanos для уникальности
|
||||
rawName := fmt.Sprintf("%s-%d", name, time.Now().UnixNano())
|
||||
jobName := rawName
|
||||
if len(jobName) > 63 {
|
||||
jobName = jobName[:63]
|
||||
}
|
||||
|
||||
fj := &slessv1alpha1.FunctionJob{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: jobName,
|
||||
Namespace: ns,
|
||||
},
|
||||
Spec: slessv1alpha1.FunctionJobSpec{
|
||||
FunctionRef: name,
|
||||
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),
|
||||
))
|
||||
}
|
||||
|
||||
@@ -0,0 +1,305 @@
|
||||
// Создано: 2026-03-20 (function-service-split)
|
||||
// services.go — CRUD handlers для Service CRD (sless_service).
|
||||
// sless_service = long-running Deployment + URL. Каждый вызов проксируется к поду.
|
||||
// Namespace берётся из URL: /v1/namespaces/{namespace}/services/{name}
|
||||
|
||||
package handler
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"io"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"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"
|
||||
"gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/builder"
|
||||
)
|
||||
|
||||
// serviceRequest — тело запроса для создания/обновления сервиса.
|
||||
type serviceRequest struct {
|
||||
Name string `json:"name"`
|
||||
Runtime string `json:"runtime"`
|
||||
Entrypoint string `json:"entrypoint"`
|
||||
MemoryMB int32 `json:"memory_mb"`
|
||||
TimeoutSec int32 `json:"timeout_sec"`
|
||||
Env map[string]string `json:"env_vars"`
|
||||
S3Bucket string `json:"s3_bucket"`
|
||||
S3Key string `json:"s3_key"`
|
||||
}
|
||||
|
||||
// serviceResponse — ответ при чтении сервиса.
|
||||
// URL — ключевое отличие от functionResponse: заполняется оператором сразу после деплоя.
|
||||
type serviceResponse struct {
|
||||
Name string `json:"name"`
|
||||
Namespace string `json:"namespace"`
|
||||
Runtime string `json:"runtime"`
|
||||
Entrypoint string `json:"entrypoint"`
|
||||
MemoryMB int32 `json:"memory_mb"`
|
||||
TimeoutSec int32 `json:"timeout_sec"`
|
||||
Env map[string]string `json:"env_vars"`
|
||||
S3Bucket string `json:"s3_bucket"`
|
||||
S3Key string `json:"s3_key"`
|
||||
Phase slessv1alpha1.ServicePhase `json:"phase"`
|
||||
ImageRef string `json:"image_ref"`
|
||||
URL string `json:"url,omitempty"`
|
||||
Message string `json:"message,omitempty"`
|
||||
CreatedAt string `json:"created_at,omitempty"`
|
||||
LastBuiltAt string `json:"last_built_at,omitempty"`
|
||||
}
|
||||
|
||||
// svcToResponse конвертирует Service CRD в ответ API.
|
||||
func svcToResponse(svc *slessv1alpha1.Service) serviceResponse {
|
||||
resp := serviceResponse{
|
||||
Name: svc.Name,
|
||||
Namespace: svc.Namespace,
|
||||
Runtime: svc.Spec.Runtime,
|
||||
Entrypoint: svc.Spec.Entrypoint,
|
||||
MemoryMB: svc.Spec.MemoryMB,
|
||||
TimeoutSec: svc.Spec.TimeoutSec,
|
||||
Env: svc.Spec.Env,
|
||||
S3Bucket: svc.Spec.S3Bucket,
|
||||
S3Key: svc.Spec.S3Key,
|
||||
Phase: svc.Status.Phase,
|
||||
ImageRef: svc.Status.ImageRef,
|
||||
URL: svc.Status.URL,
|
||||
Message: svc.Status.Message,
|
||||
}
|
||||
if !svc.CreationTimestamp.IsZero() {
|
||||
resp.CreatedAt = svc.CreationTimestamp.UTC().Format("2006-01-02 15:04:05 UTC")
|
||||
}
|
||||
if svc.Status.LastBuiltAt != nil && !svc.Status.LastBuiltAt.IsZero() {
|
||||
resp.LastBuiltAt = svc.Status.LastBuiltAt.UTC().Format("2006-01-02 15:04:05 UTC")
|
||||
}
|
||||
return resp
|
||||
}
|
||||
|
||||
// ListServices — GET /v1/namespaces/{namespace}/services
|
||||
func (h *Handler) ListServices(w http.ResponseWriter, r *http.Request) {
|
||||
ns := namespace(r)
|
||||
list := &slessv1alpha1.ServiceList{}
|
||||
if err := h.K8s.List(r.Context(), list, client.InNamespace(ns)); err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||||
return
|
||||
}
|
||||
result := make([]serviceResponse, 0, len(list.Items))
|
||||
for i := range list.Items {
|
||||
result = append(result, svcToResponse(&list.Items[i]))
|
||||
}
|
||||
writeJSON(w, http.StatusOK, result)
|
||||
}
|
||||
|
||||
// CreateService — POST /v1/namespaces/{namespace}/services
|
||||
func (h *Handler) CreateService(w http.ResponseWriter, r *http.Request) {
|
||||
ns := namespace(r)
|
||||
var req serviceRequest
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeJSON(w, http.StatusBadRequest, errResp("invalid JSON: "+err.Error()))
|
||||
return
|
||||
}
|
||||
if req.Name == "" || req.Runtime == "" {
|
||||
writeJSON(w, http.StatusBadRequest, errResp("name and runtime are required"))
|
||||
return
|
||||
}
|
||||
if req.Entrypoint == "" {
|
||||
writeJSON(w, http.StatusBadRequest, errResp("entrypoint is required"))
|
||||
return
|
||||
}
|
||||
if req.MemoryMB <= 0 || req.MemoryMB > 4096 {
|
||||
writeJSON(w, http.StatusBadRequest, errResp("memory_mb must be between 1 and 4096"))
|
||||
return
|
||||
}
|
||||
|
||||
svc := &slessv1alpha1.Service{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: req.Name,
|
||||
Namespace: ns,
|
||||
},
|
||||
Spec: slessv1alpha1.ServiceSpec{
|
||||
Runtime: req.Runtime,
|
||||
Entrypoint: req.Entrypoint,
|
||||
MemoryMB: req.MemoryMB,
|
||||
TimeoutSec: req.TimeoutSec,
|
||||
Env: req.Env,
|
||||
S3Bucket: req.S3Bucket,
|
||||
S3Key: req.S3Key,
|
||||
},
|
||||
}
|
||||
if err := h.K8s.Create(r.Context(), svc); err != nil {
|
||||
if errors.IsAlreadyExists(err) {
|
||||
// Повторяем логику function_handler: обрабатываем split-brain кеша.
|
||||
// Если объект реально есть и не в фазе Failed — возвращаем 409.
|
||||
existing := &slessv1alpha1.Service{}
|
||||
getErr := h.K8s.Get(r.Context(), client.ObjectKey{Name: req.Name, Namespace: ns}, existing)
|
||||
shouldRecreate := errors.IsNotFound(getErr) ||
|
||||
(getErr == nil && existing.Status.Phase == slessv1alpha1.ServicePhaseFailed)
|
||||
if shouldRecreate {
|
||||
if getErr == nil {
|
||||
_ = h.K8s.Delete(r.Context(), existing)
|
||||
}
|
||||
svc.ResourceVersion = ""
|
||||
if createErr := h.K8s.Create(r.Context(), svc); createErr != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, errResp(createErr.Error()))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusCreated, svcToResponse(svc))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusConflict, errResp("service already exists"))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusCreated, svcToResponse(svc))
|
||||
}
|
||||
|
||||
// GetService — GET /v1/namespaces/{namespace}/services/{name}
|
||||
func (h *Handler) GetService(w http.ResponseWriter, r *http.Request) {
|
||||
ns := namespace(r)
|
||||
name := pathVar(r, "name")
|
||||
svc := &slessv1alpha1.Service{}
|
||||
if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, svc); err != nil {
|
||||
if errors.IsNotFound(err) {
|
||||
writeJSON(w, http.StatusNotFound, errResp("service not found"))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, svcToResponse(svc))
|
||||
}
|
||||
|
||||
// UpdateService — PUT /v1/namespaces/{namespace}/services/{name}
|
||||
func (h *Handler) UpdateService(w http.ResponseWriter, r *http.Request) {
|
||||
ns := namespace(r)
|
||||
name := pathVar(r, "name")
|
||||
var req serviceRequest
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeJSON(w, http.StatusBadRequest, errResp("invalid JSON: "+err.Error()))
|
||||
return
|
||||
}
|
||||
|
||||
svc := &slessv1alpha1.Service{}
|
||||
if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, svc); err != nil {
|
||||
if errors.IsNotFound(err) {
|
||||
writeJSON(w, http.StatusNotFound, errResp("service not found"))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||||
return
|
||||
}
|
||||
|
||||
if req.Runtime == "" || req.Entrypoint == "" {
|
||||
writeJSON(w, http.StatusBadRequest, errResp("runtime and entrypoint are required"))
|
||||
return
|
||||
}
|
||||
if req.MemoryMB <= 0 || req.MemoryMB > 4096 {
|
||||
writeJSON(w, http.StatusBadRequest, errResp("memory_mb must be between 1 and 4096"))
|
||||
return
|
||||
}
|
||||
|
||||
svc.Spec.Runtime = req.Runtime
|
||||
svc.Spec.Entrypoint = req.Entrypoint
|
||||
svc.Spec.MemoryMB = req.MemoryMB
|
||||
svc.Spec.TimeoutSec = req.TimeoutSec
|
||||
svc.Spec.Env = req.Env
|
||||
if req.S3Bucket != "" {
|
||||
svc.Spec.S3Bucket = req.S3Bucket
|
||||
}
|
||||
if req.S3Key != "" {
|
||||
svc.Spec.S3Key = req.S3Key
|
||||
}
|
||||
|
||||
if err := h.K8s.Update(r.Context(), svc); err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, svcToResponse(svc))
|
||||
}
|
||||
|
||||
// DeleteService — DELETE /v1/namespaces/{namespace}/services/{name}
|
||||
func (h *Handler) DeleteService(w http.ResponseWriter, r *http.Request) {
|
||||
ns := namespace(r)
|
||||
name := pathVar(r, "name")
|
||||
svc := &slessv1alpha1.Service{}
|
||||
if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, svc); err != nil {
|
||||
if errors.IsNotFound(err) {
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||||
return
|
||||
}
|
||||
if err := h.K8s.Delete(r.Context(), svc); err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
}
|
||||
|
||||
// UploadServiceCode — POST /v1/namespaces/{namespace}/services/{name}/upload
|
||||
// Принимает multipart/form-data с полем "code" (zip архив с кодом сервиса).
|
||||
// Обновляет Service CRD — оператор запустит kaniko и затем задеплоит Deployment.
|
||||
// Логика идентична UploadCode (functions), но работает с Service CRD.
|
||||
func (h *Handler) UploadServiceCode(w http.ResponseWriter, r *http.Request) {
|
||||
ns := namespace(r)
|
||||
name := pathVar(r, "name")
|
||||
|
||||
svc := &slessv1alpha1.Service{}
|
||||
if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, svc); err != nil {
|
||||
if errors.IsNotFound(err) {
|
||||
writeJSON(w, http.StatusNotFound, errResp("service not found"))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||||
return
|
||||
}
|
||||
|
||||
if err := r.ParseMultipartForm(32 << 20); err != nil {
|
||||
writeJSON(w, http.StatusBadRequest, errResp("invalid multipart form: "+err.Error()))
|
||||
return
|
||||
}
|
||||
file, _, err := r.FormFile("code")
|
||||
if err != nil {
|
||||
writeJSON(w, http.StatusBadRequest, errResp(`field "code" is required (zip file)`))
|
||||
return
|
||||
}
|
||||
defer file.Close()
|
||||
|
||||
zipData, err := io.ReadAll(file)
|
||||
if err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, errResp("read upload: "+err.Error()))
|
||||
return
|
||||
}
|
||||
|
||||
buf, err := builder.PrepareContext(zipData, svc.Spec.Runtime)
|
||||
if err != nil {
|
||||
writeJSON(w, http.StatusBadRequest, errResp("prepare build context: "+err.Error()))
|
||||
return
|
||||
}
|
||||
|
||||
version := time.Now().Format("20060102150405")
|
||||
s3Key, err := h.S3.UploadContext(r.Context(), ns, name, version, buf, int64(buf.Len()))
|
||||
if err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, errResp("upload to S3: "+err.Error()))
|
||||
return
|
||||
}
|
||||
|
||||
patch := client.MergeFrom(svc.DeepCopy())
|
||||
svc.Spec.S3Key = s3Key
|
||||
svc.Spec.S3Bucket = h.S3.Bucket()
|
||||
if err := h.K8s.Patch(r.Context(), svc, patch); err != nil {
|
||||
writeJSON(w, http.StatusInternalServerError, errResp("update service: "+err.Error()))
|
||||
return
|
||||
}
|
||||
|
||||
writeJSON(w, http.StatusOK, map[string]string{
|
||||
"s3_key": s3Key,
|
||||
"phase": string(slessv1alpha1.ServicePhasePending),
|
||||
"message": "build queued",
|
||||
})
|
||||
}
|
||||
@@ -1,4 +1,4 @@
|
||||
// Изменено: 2026-04-25
|
||||
// Изменено: 2026-03-20 (function-service-split: добавлены /services маршруты)
|
||||
// router.go — регистрация всех REST-маршрутов через gorilla/mux.
|
||||
// Все маршруты защищены Bearer-токеном (middleware.Auth).
|
||||
// Маршруты сгруппированы по /v1/namespaces/{namespace}/...
|
||||
@@ -48,6 +48,14 @@ func NewRouter(h *handler.Handler, log *slog.Logger) http.Handler {
|
||||
// Source code — возвращает файлы из tar.gz контекста сборки (без Dockerfile)
|
||||
v1.HandleFunc("/namespaces/{namespace}/functions/{name}/source", h.GetSource).Methods(http.MethodGet)
|
||||
|
||||
// Services CRUD — long-running Deployment + URL (sless_service)
|
||||
v1.HandleFunc("/namespaces/{namespace}/services", h.ListServices).Methods(http.MethodGet)
|
||||
v1.HandleFunc("/namespaces/{namespace}/services", h.CreateService).Methods(http.MethodPost)
|
||||
v1.HandleFunc("/namespaces/{namespace}/services/{name}", h.GetService).Methods(http.MethodGet)
|
||||
v1.HandleFunc("/namespaces/{namespace}/services/{name}", h.UpdateService).Methods(http.MethodPut)
|
||||
v1.HandleFunc("/namespaces/{namespace}/services/{name}", h.DeleteService).Methods(http.MethodDelete)
|
||||
v1.HandleFunc("/namespaces/{namespace}/services/{name}/upload", h.UploadServiceCode).Methods(http.MethodPost)
|
||||
|
||||
// Triggers CRUD
|
||||
v1.HandleFunc("/namespaces/{namespace}/triggers", h.ListTriggers).Methods(http.MethodGet)
|
||||
v1.HandleFunc("/namespaces/{namespace}/triggers", h.CreateTrigger).Methods(http.MethodPost)
|
||||
|
||||
Reference in New Issue
Block a user