Files
sless/internal/api/handler/jobs.go
T
“Naeel” f41cd39b26 feat: namespace-per-user via JWT sub SHA256 + ensureNamespace in operator
- operator: ensureNamespace() создаёт k8s namespace при первом Create-запросе
- operator: defaultNamespace константа вместо хардкода 'default'
- provider: SubFromJWT декодирует JWT payload, извлекает sub
- provider: NamespaceFromSub вычисляет sless-{sha256[:8]} из sub
- provider: PingNubesAPI валидирует токен запросом к nubes API
- provider: Configure вычисляет namespace и создаёт Client с ним
- provider: новый атрибут nubes_endpoint (опционально, env: NUBES_ENDPOINT)
2026-03-11 07:35:49 +04:00

160 lines
5.1 KiB
Go

// Изменено: 2026-03-08
// jobs.go — CRUD handlers для FunctionJob CRD.
// Создаёт/читает/удаляет k8s FunctionJob ресурсы.
// Namespace берётся из URL: /v1/namespaces/{namespace}/jobs/{name}
package handler
import (
"encoding/json"
"net/http"
"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"
)
// jobRequest — тело POST /v1/namespaces/{ns}/jobs
type jobRequest struct {
Name string `json:"name"`
FunctionRef string `json:"function"`
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"`
FunctionRef string `json:"function"`
EventJSON string `json:"event_json"`
RunID int64 `json:"run_id"`
Phase string `json:"phase"`
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,
FunctionRef: j.Spec.FunctionRef,
EventJSON: j.Spec.EventJSON,
RunID: j.Spec.RunID,
Phase: string(j.Status.Phase),
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. Оператор запустит k8s Job асинхронно.
func (h *Handler) CreateJob(w http.ResponseWriter, r *http.Request) {
ns := namespace(r)
// Гарантируем существование namespace — аналогично CreateFunction.
if err := h.ensureNamespace(r.Context(), ns); err != nil {
writeJSON(w, http.StatusInternalServerError, errResp("ensure namespace: "+err.Error()))
return
}
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.FunctionRef == "" {
writeJSON(w, http.StatusBadRequest, errResp("function is required"))
return
}
if req.EventJSON == "" {
req.EventJSON = "{}"
}
job := &slessv1alpha1.FunctionJob{
ObjectMeta: metav1.ObjectMeta{
Name: req.Name,
Namespace: ns,
},
Spec: slessv1alpha1.FunctionJobSpec{
FunctionRef: req.FunctionRef,
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) {
w.WriteHeader(http.StatusNoContent)
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)
}