Files
sless/internal/api/handler/services.go
T
Naeel a76baa62a3 fix: v0.1.50 — runtime 400 + SLESS_ENTRYPOINT в Deployment
Bug 1: services.go — k8s IsInvalid error (CRD enum validation) маппился в 500.
Теперь errors.IsInvalid() → 400 Bad Request (invalid service spec).

Bug 2: service_controller.go buildServiceDeployment не передавал env SLESS_ENTRYPOINT
в под. Добавлен в envVars из svc.Spec.Entrypoint. Без него server.py использовал
fallback handler.handle и не замечал неверный entrypoint.

operator_failure_test.sh 12B-2: обновлён под новое правильное поведение —
create ruby3.0 → 400 (не 201). Старый 201-путь сохранён как warn для совместимости.
2026-03-22 08:11:05 +03:00

323 lines
12 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Создано: 2026-03-20 (function-service-split)
// Изменено: 2026-03-21 (fix: DeleteService возвращает 404 вместо 204 при отсутствующем объекте)
// 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
}
// timeout_sec: 0 = нет ограничения (допустимо), отрицательное — ошибка.
if req.TimeoutSec < 0 || req.TimeoutSec > 900 {
writeJSON(w, http.StatusBadRequest, errResp("timeout_sec must be 0 (no limit) or 1900"))
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
}
// k8s возвращает StatusInvalid (422) при нарушении enum-валидации CRD (например, неизвестный runtime)
// — маппим это в 400, а не в 500
if errors.IsInvalid(err) {
writeJSON(w, http.StatusBadRequest, errResp("invalid service spec: "+err.Error()))
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
}
// timeout_sec: 0 = нет ограничения (допустимо), отрицательное — ошибка.
if req.TimeoutSec < 0 || req.TimeoutSec > 900 {
writeJSON(w, http.StatusBadRequest, errResp("timeout_sec must be 0 (no limit) or 1900"))
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) {
writeJSON(w, http.StatusNotFound, errResp("service not found"))
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",
})
}