- handler.go: убраны бизнес-логика и k8s-типы (corev1/k8serrors/metav1)
handler.go теперь только инфраструктура: Handler struct + helpers
- namespace.go: новый файл — EnsureNamespace хендлер живёт здесь
SoC: создание namespace — отдельная ответственность, не смешивается с CRUD
- router.go: добавлен маршрут POST /v1/namespaces/{namespace}/ensure
- client.go: добавлен метод EnsureNamespace(ctx, ns) → POST /ensure
- provider.go: Configure() вызывает c.EnsureNamespace(ctx, namespace) после создания Client
Namespace создаётся ОДИН РАЗ при инициализации провайдера
Resource-хендлеры (Function, Trigger, Job) namespace не трогают
- .gitignore: добавлена директория secrets/ (токены, ключи)
- provider v0.1.13, operator v0.1.21
Operator: naeel/sless-operator:v0.1.21
Provider: terra.k8c.ru/naeel/sless v0.1.13
155 lines
4.9 KiB
Go
155 lines
4.9 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)
|
|
|
|
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)
|
|
}
|