262 lines
9.2 KiB
Go
262 lines
9.2 KiB
Go
// Изменено: 2026-03-20 (merge: FunctionJob теперь самодостаточен — убран FunctionRef, добавлены Runtime/Entrypoint/Env)
|
||
// Изменено: 2026-03-21 (fix: DeleteJob возвращает 404 вместо 204 при отсутствующем объекте)
|
||
// Изменено: 2026-03-22 (fix: UploadJobCode retry loop против cache lag controller-runtime)
|
||
// jobs.go — CRUD handlers для FunctionJob CRD.
|
||
// Создаёт/читает/удаляет k8s FunctionJob ресурсы.
|
||
// Namespace берётся из URL: /v1/namespaces/{namespace}/jobs/{name}
|
||
|
||
package handler
|
||
|
||
import (
|
||
"crypto/sha256"
|
||
"encoding/json"
|
||
"fmt"
|
||
"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"
|
||
)
|
||
|
||
// jobRequest — тело POST /v1/namespaces/{ns}/jobs
|
||
type jobRequest struct {
|
||
Name string `json:"name"`
|
||
Runtime string `json:"runtime"`
|
||
Entrypoint string `json:"entrypoint"`
|
||
MemoryMB int32 `json:"memory_mb,omitempty"`
|
||
TimeoutSec int32 `json:"timeout_sec,omitempty"`
|
||
Env map[string]string `json:"env_vars,omitempty"`
|
||
S3Bucket string `json:"s3_bucket,omitempty"`
|
||
S3Key string `json:"s3_key,omitempty"`
|
||
EventJSON string `json:"event_json,omitempty"`
|
||
// RunID — идентификатор запуска. 0 = создать без запуска, >0 = запустить.
|
||
RunID int64 `json:"run_id"`
|
||
}
|
||
|
||
// jobResponse — ответ при чтении / создании FunctionJob
|
||
type jobResponse struct {
|
||
Name string `json:"name"`
|
||
Namespace string `json:"namespace"`
|
||
Runtime string `json:"runtime"`
|
||
Entrypoint string `json:"entrypoint"`
|
||
EventJSON string `json:"event_json"`
|
||
RunID int64 `json:"run_id"`
|
||
Phase string `json:"phase"`
|
||
ImageRef string `json:"image_ref,omitempty"`
|
||
JobName string `json:"job_name,omitempty"`
|
||
StartTime string `json:"start_time,omitempty"`
|
||
CompletionTime string `json:"completion_time,omitempty"`
|
||
Message string `json:"message,omitempty"`
|
||
}
|
||
|
||
// jobToResponse конвертирует FunctionJob CR → jobResponse.
|
||
func jobToResponse(j *slessv1alpha1.FunctionJob) jobResponse {
|
||
r := jobResponse{
|
||
Name: j.Name,
|
||
Namespace: j.Namespace,
|
||
Runtime: j.Spec.Runtime,
|
||
Entrypoint: j.Spec.Entrypoint,
|
||
EventJSON: j.Spec.EventJSON,
|
||
RunID: j.Spec.RunID,
|
||
Phase: string(j.Status.Phase),
|
||
ImageRef: j.Status.ImageRef,
|
||
JobName: j.Status.JobName,
|
||
Message: j.Status.Message,
|
||
}
|
||
if j.Status.StartTime != nil {
|
||
r.StartTime = j.Status.StartTime.UTC().Format("2006-01-02T15:04:05Z")
|
||
}
|
||
if j.Status.CompletionTime != nil {
|
||
r.CompletionTime = j.Status.CompletionTime.UTC().Format("2006-01-02T15:04:05Z")
|
||
}
|
||
return r
|
||
}
|
||
|
||
// CreateJob — POST /v1/namespaces/{namespace}/jobs
|
||
// Создаёт FunctionJob CR. Оператор запустит kaniko сборку и затем k8s Job асинхронно.
|
||
func (h *Handler) CreateJob(w http.ResponseWriter, r *http.Request) {
|
||
ns := namespace(r)
|
||
|
||
var req jobRequest
|
||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||
writeJSON(w, http.StatusBadRequest, errResp("invalid JSON: "+err.Error()))
|
||
return
|
||
}
|
||
if req.Name == "" {
|
||
writeJSON(w, http.StatusBadRequest, errResp("name is required"))
|
||
return
|
||
}
|
||
if req.Runtime == "" {
|
||
writeJSON(w, http.StatusBadRequest, errResp("runtime is required"))
|
||
return
|
||
}
|
||
if req.Entrypoint == "" {
|
||
writeJSON(w, http.StatusBadRequest, errResp("entrypoint is required"))
|
||
return
|
||
}
|
||
if req.EventJSON == "" {
|
||
req.EventJSON = "{}"
|
||
}
|
||
|
||
job := &slessv1alpha1.FunctionJob{
|
||
ObjectMeta: metav1.ObjectMeta{
|
||
Name: req.Name,
|
||
Namespace: ns,
|
||
},
|
||
Spec: slessv1alpha1.FunctionJobSpec{
|
||
Runtime: req.Runtime,
|
||
Entrypoint: req.Entrypoint,
|
||
MemoryMB: req.MemoryMB,
|
||
TimeoutSec: req.TimeoutSec,
|
||
Env: req.Env,
|
||
S3Bucket: req.S3Bucket,
|
||
S3Key: req.S3Key,
|
||
EventJSON: req.EventJSON,
|
||
RunID: req.RunID,
|
||
},
|
||
}
|
||
|
||
if err := h.K8s.Create(r.Context(), job); err != nil {
|
||
if errors.IsAlreadyExists(err) {
|
||
writeJSON(w, http.StatusConflict, errResp("job already exists"))
|
||
return
|
||
}
|
||
h.Log.Error("create FunctionJob", "err", err)
|
||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||
return
|
||
}
|
||
|
||
writeJSON(w, http.StatusCreated, jobToResponse(job))
|
||
}
|
||
|
||
// GetJob — GET /v1/namespaces/{namespace}/jobs/{name}
|
||
// Возвращает статус FunctionJob включая phase и время выполнения.
|
||
func (h *Handler) GetJob(w http.ResponseWriter, r *http.Request) {
|
||
ns := namespace(r)
|
||
name := pathVar(r, "name")
|
||
|
||
var job slessv1alpha1.FunctionJob
|
||
if err := h.K8s.Get(r.Context(), client.ObjectKey{Namespace: ns, Name: name}, &job); err != nil {
|
||
if errors.IsNotFound(err) {
|
||
writeJSON(w, http.StatusNotFound, errResp("job not found"))
|
||
return
|
||
}
|
||
h.Log.Error("get FunctionJob", "err", err)
|
||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||
return
|
||
}
|
||
|
||
writeJSON(w, http.StatusOK, jobToResponse(&job))
|
||
}
|
||
|
||
// DeleteJob — DELETE /v1/namespaces/{namespace}/jobs/{name}
|
||
// Удаляет FunctionJob CR (и подчинённый k8s Job через ownerReference, если ещё не убран по TTL).
|
||
func (h *Handler) DeleteJob(w http.ResponseWriter, r *http.Request) {
|
||
ns := namespace(r)
|
||
name := pathVar(r, "name")
|
||
|
||
var job slessv1alpha1.FunctionJob
|
||
if err := h.K8s.Get(r.Context(), client.ObjectKey{Namespace: ns, Name: name}, &job); err != nil {
|
||
if errors.IsNotFound(err) {
|
||
writeJSON(w, http.StatusNotFound, errResp("job not found"))
|
||
return
|
||
}
|
||
h.Log.Error("get FunctionJob for delete", "err", err)
|
||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||
return
|
||
}
|
||
|
||
if err := h.K8s.Delete(r.Context(), &job); err != nil {
|
||
h.Log.Error("delete FunctionJob", "err", err)
|
||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||
return
|
||
}
|
||
|
||
w.WriteHeader(http.StatusNoContent)
|
||
}
|
||
|
||
// UploadJobCode — POST /v1/namespaces/{namespace}/jobs/{name}/upload
|
||
// Принимает multipart/form-data с полем "code" (zip архив с кодом функции).
|
||
// Аналогично UploadCode для Function, но работает с FunctionJob CRD.
|
||
// После загрузки обновляет Spec.S3Key — контроллер начнёт kaniko сборку.
|
||
func (h *Handler) UploadJobCode(w http.ResponseWriter, r *http.Request) {
|
||
ns := namespace(r)
|
||
name := pathVar(r, "name")
|
||
|
||
// Читаем FunctionJob для получения runtime.
|
||
// Retry до 5 раз с задержкой 200мс — защита от cache lag controller-runtime:
|
||
// сразу после POST /jobs (201) informer cache может ещё не синхронизировать новый CR.
|
||
fj := &slessv1alpha1.FunctionJob{}
|
||
var getJobErr error
|
||
for i := 0; i < 5; i++ {
|
||
getJobErr = h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, fj)
|
||
if getJobErr == nil || !errors.IsNotFound(getJobErr) {
|
||
break
|
||
}
|
||
time.Sleep(200 * time.Millisecond)
|
||
}
|
||
if getJobErr != nil {
|
||
if errors.IsNotFound(getJobErr) {
|
||
writeJSON(w, http.StatusNotFound, errResp("job not found"))
|
||
return
|
||
}
|
||
writeJSON(w, http.StatusInternalServerError, errResp(getJobErr.Error()))
|
||
return
|
||
}
|
||
|
||
// Лимит 32MB на загрузку кода функции
|
||
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
|
||
}
|
||
|
||
// Готовим build context: Dockerfile + tar.gz для kaniko
|
||
buf, err := builder.PrepareContext(zipData, fj.Spec.Runtime)
|
||
if err != nil {
|
||
writeJSON(w, http.StatusBadRequest, errResp("prepare build context: "+err.Error()))
|
||
return
|
||
}
|
||
|
||
// Версия = sha256(загруженный zip). Одинаковый код → одинаковый s3Key → cache hit.
|
||
// Изменился код → новый хеш → новый build.
|
||
zipHash := sha256.Sum256(zipData)
|
||
version := fmt.Sprintf("%x", zipHash[:])[:16]
|
||
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
|
||
}
|
||
|
||
// Обновляем FunctionJob CRD: новый s3Key → контроллер начнёт сборку
|
||
patch := client.MergeFrom(fj.DeepCopy())
|
||
fj.Spec.S3Key = s3Key
|
||
fj.Spec.S3Bucket = h.S3.Bucket()
|
||
if err := h.K8s.Patch(r.Context(), fj, patch); err != nil {
|
||
writeJSON(w, http.StatusInternalServerError, errResp("update job: "+err.Error()))
|
||
return
|
||
}
|
||
|
||
writeJSON(w, http.StatusOK, map[string]string{
|
||
"s3_key": s3Key,
|
||
"phase": string(slessv1alpha1.FunctionJobPhasePending),
|
||
"message": "build queued",
|
||
})
|
||
}
|