Files
sless/internal/api/handler/services.go
T
Naeel b86ff3a62e feat(service): timeout_sec без дефолта; 0=нет таймаута; operator v0.1.48
- 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
2026-03-21 16:58:43 +03:00

317 lines
11 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
}
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",
})
}