feat(job): merge sless_function into sless_job — self-contained build+run
- FunctionJobSpec: убран FunctionRef, добавлены Runtime/Entrypoint/Env/S3Key/MemoryMB/TimeoutSec
- FunctionJobStatus: новый ImageRef, новая фаза Building
- FunctionJobReconciler: Building фаза (kaniko), убрана зависимость от Function CRD
Builder+OperatorNamespace как поля struct; аннотация sless.kube5s.ru/build-job guard
- main.go: Builder+OperatorNamespace переданы в FunctionJobReconciler
- jobs.go handler: jobRequest/jobResponse без FunctionRef; новый UploadJobCode handler
- router.go: /jobs/{name}/upload маршрут
- client.go: JobRequest/JobResponse обновлены; UploadJobCode; uploadCodeToURL общий хелпер
- job_resource.go: полная переработка — источник/среда встроены в JobModel, ModifyPlan,
Create с upload, wait_timeout_sec=900 по умолчанию (kaniko + выполнение)
- examples/POSTGRES/functions.tf: раскомментирован, sless_function удалён,
sless_job самодостаточен (inline source_dir/runtime/entrypoint/env_vars)
This commit is contained in:
@@ -65,3 +65,7 @@ plan.out
|
||||
sless-plan
|
||||
examples/.git
|
||||
event-dispatcher
|
||||
|
||||
# build artifacts
|
||||
/sless
|
||||
examples/POSTGRES/stress_log*.txt
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
// Изменено: 2026-03-08
|
||||
// Изменено: 2026-03-20 (merge sless_function+sless_job: FunctionJobSpec самодостаточен,
|
||||
// больше не требует отдельного Function CRD)
|
||||
// Описание CRD FunctionJob — одноразовый запуск функции.
|
||||
// Отдельный ресурс (не Trigger) потому что семантика принципиально другая:
|
||||
// - Trigger: постоянно живёт, описывает КАК функцию вызывают (http/cron)
|
||||
@@ -14,6 +15,8 @@ import (
|
||||
)
|
||||
|
||||
// FunctionJobSpec — параметры одноразового запуска функции.
|
||||
// Самодостаточен: содержит всё для сборки образа и запуска Job.
|
||||
// Отдельный Function CRD больше не требуется.
|
||||
type FunctionJobSpec struct {
|
||||
// RunID — идентификатор запуска. 0 = не запускать.
|
||||
// Каждое ненулевое значение уникально идентифицирует запуск.
|
||||
@@ -22,9 +25,35 @@ type FunctionJobSpec struct {
|
||||
// +kubebuilder:default=0
|
||||
RunID int64 `json:"runId"`
|
||||
|
||||
// FunctionRef — имя Function ресурса в том же namespace
|
||||
// --- Параметры функции (встроены, не нужен отдельный sless_function) ---
|
||||
|
||||
// Runtime — язык и версия выполнения (go1.23, python3.11, nodejs20)
|
||||
// +kubebuilder:validation:Enum=go1.23;python3.11;nodejs20
|
||||
// +kubebuilder:validation:Required
|
||||
FunctionRef string `json:"functionRef"`
|
||||
Runtime string `json:"runtime"`
|
||||
|
||||
// Entrypoint — точка входа в код функции (например: handler.handle)
|
||||
// +kubebuilder:validation:Required
|
||||
Entrypoint string `json:"entrypoint"`
|
||||
|
||||
// S3Bucket — бакет S3 где хранится tar.gz контекст сборки
|
||||
S3Bucket string `json:"s3Bucket,omitempty"`
|
||||
|
||||
// S3Key — ключ объекта в S3 (путь до tar.gz контекста)
|
||||
S3Key string `json:"s3Key,omitempty"`
|
||||
|
||||
// MemoryMB — лимит памяти в мегабайтах (default: 128)
|
||||
// +kubebuilder:default=128
|
||||
MemoryMB int32 `json:"memoryMB,omitempty"`
|
||||
|
||||
// TimeoutSec — максимальное время выполнения в секундах (default: 30)
|
||||
// +kubebuilder:default=30
|
||||
TimeoutSec int32 `json:"timeoutSec,omitempty"`
|
||||
|
||||
// Env — переменные окружения, передаются в контейнер функции
|
||||
Env map[string]string `json:"env,omitempty"`
|
||||
|
||||
// --- Job-специфичные параметры ---
|
||||
|
||||
// EventJSON — данные передаваемые в handle(event) в JSON формате.
|
||||
// Если не задан — передаётся пустой объект {}.
|
||||
@@ -37,8 +66,10 @@ type FunctionJobSpec struct {
|
||||
type FunctionJobPhase string
|
||||
|
||||
const (
|
||||
// FunctionJobPhasePending — ожидает пока Function станет Ready
|
||||
// FunctionJobPhasePending — ожидает начала сборки или запуска
|
||||
FunctionJobPhasePending FunctionJobPhase = "Pending"
|
||||
// FunctionJobPhaseBuilding — идёт сборка Docker образа через kaniko
|
||||
FunctionJobPhaseBuilding FunctionJobPhase = "Building"
|
||||
// FunctionJobPhaseRunning — k8s Job запущен, функция выполняется
|
||||
FunctionJobPhaseRunning FunctionJobPhase = "Running"
|
||||
// FunctionJobPhaseSucceeded — функция успешно завершилась
|
||||
@@ -49,12 +80,16 @@ const (
|
||||
|
||||
// FunctionJobStatus — наблюдаемое состояние (заполняет контроллер).
|
||||
type FunctionJobStatus struct {
|
||||
// Phase — текущая фаза: Pending, Running, Succeeded, Failed
|
||||
// Phase — текущая фаза: Pending, Building, Running, Succeeded, Failed
|
||||
Phase FunctionJobPhase `json:"phase,omitempty"`
|
||||
|
||||
// JobName — имя созданного k8s Job
|
||||
JobName string `json:"jobName,omitempty"`
|
||||
|
||||
// ImageRef — полный путь к собранному Docker образу в registry
|
||||
// Заполняется после успешной сборки (фаза Building → Running).
|
||||
ImageRef string `json:"imageRef,omitempty"`
|
||||
|
||||
// StartTime — время запуска k8s Job
|
||||
StartTime *metav1.Time `json:"startTime,omitempty"`
|
||||
|
||||
@@ -67,7 +102,7 @@ type FunctionJobStatus struct {
|
||||
|
||||
//+kubebuilder:object:root=true
|
||||
//+kubebuilder:subresource:status
|
||||
//+kubebuilder:printcolumn:name="Function",type=string,JSONPath=`.spec.functionRef`
|
||||
//+kubebuilder:printcolumn:name="Runtime",type=string,JSONPath=`.spec.runtime`
|
||||
//+kubebuilder:printcolumn:name="Phase",type=string,JSONPath=`.status.phase`
|
||||
//+kubebuilder:printcolumn:name="Age",type=date,JSONPath=`.metadata.creationTimestamp`
|
||||
|
||||
|
||||
@@ -57,7 +57,7 @@ func (in *FunctionJob) DeepCopyInto(out *FunctionJob) {
|
||||
*out = *in
|
||||
out.TypeMeta = in.TypeMeta
|
||||
in.ObjectMeta.DeepCopyInto(&out.ObjectMeta)
|
||||
out.Spec = in.Spec
|
||||
in.Spec.DeepCopyInto(&out.Spec)
|
||||
in.Status.DeepCopyInto(&out.Status)
|
||||
}
|
||||
|
||||
@@ -112,8 +112,16 @@ func (in *FunctionJobList) DeepCopyObject() runtime.Object {
|
||||
}
|
||||
|
||||
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
|
||||
// FunctionJobSpec содержит map[string]string Env — требует явного deep copy.
|
||||
func (in *FunctionJobSpec) DeepCopyInto(out *FunctionJobSpec) {
|
||||
*out = *in
|
||||
if in.Env != nil {
|
||||
in, out := &in.Env, &out.Env
|
||||
*out = make(map[string]string, len(*in))
|
||||
for key, val := range *in {
|
||||
(*out)[key] = val
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new FunctionJobSpec.
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
// Изменено: 2026-03-17 20:00 (bugfix: job-name label удалён в k8s 1.27+, split-brain cached client)
|
||||
// Изменено: 2026-03-20 (merge sless_function+sless_job: FunctionJobReconciler самодостаточен)
|
||||
// FunctionJobReconciler — контроллер одноразовых запусков функций.
|
||||
// При создании FunctionJob:
|
||||
// 1. Ждёт пока Function станет Ready
|
||||
// 2. Создаёт k8s Job который запускает образ функции с CMD runner
|
||||
// При создании FunctionJob с RunID>0:
|
||||
// 1. Запускает kaniko сборку образа (фаза Building) — больше не зависит от Function CRD
|
||||
// 2. После сборки создаёт k8s Job который запускает образ функции с CMD runner
|
||||
// 3. Следит за завершением Job → обновляет статус (Succeeded/Failed)
|
||||
//
|
||||
// Почему отдельный ресурс (не Trigger type=job):
|
||||
@@ -31,14 +31,17 @@ import (
|
||||
"sigs.k8s.io/controller-runtime/pkg/log"
|
||||
|
||||
slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1"
|
||||
"gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/builder"
|
||||
)
|
||||
|
||||
// FunctionJobReconciler reconciles a FunctionJob object
|
||||
type FunctionJobReconciler struct {
|
||||
client.Client
|
||||
Scheme *runtime.Scheme
|
||||
RegistrySecret string // имя k8s Secret с docker credentials (для imagePullSecrets)
|
||||
KubeClient kubernetes.Interface // typed client для чтения логов подов (logs API недоступен через controller-runtime client)
|
||||
Scheme *runtime.Scheme
|
||||
RegistrySecret string // имя k8s Secret с docker credentials (для imagePullSecrets)
|
||||
KubeClient kubernetes.Interface // typed client для чтения логов подов (logs API недоступен через controller-runtime client)
|
||||
Builder *builder.Builder // kaniko builder — собирает образ функции
|
||||
OperatorNamespace string // namespace оператора — где запускаются kaniko job'ы (обычно "sless")
|
||||
}
|
||||
|
||||
//+kubebuilder:rbac:groups=sless.kube5s.ru,resources=functionjobs,verbs=get;list;watch;create;update;patch;delete
|
||||
@@ -48,8 +51,6 @@ type FunctionJobReconciler struct {
|
||||
|
||||
// Reconcile — основной цикл контроллера.
|
||||
func (r *FunctionJobReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
|
||||
logger := log.FromContext(ctx)
|
||||
|
||||
fj := &slessv1alpha1.FunctionJob{}
|
||||
if err := r.Get(ctx, req.NamespacedName, fj); err != nil {
|
||||
if errors.IsNotFound(err) {
|
||||
@@ -76,29 +77,28 @@ func (r *FunctionJobReconciler) Reconcile(ctx context.Context, req ctrl.Request)
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
// Проверяем что Function существует и готова
|
||||
fn := &slessv1alpha1.Function{}
|
||||
if err := r.Get(ctx, client.ObjectKey{Name: fj.Spec.FunctionRef, Namespace: fj.Namespace}, fn); err != nil {
|
||||
if errors.IsNotFound(err) {
|
||||
fj.Status.Phase = slessv1alpha1.FunctionJobPhasePending
|
||||
fj.Status.Message = "function not found: " + fj.Spec.FunctionRef
|
||||
_ = r.Status().Update(ctx, fj)
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
return ctrl.Result{}, fmt.Errorf("get function: %w", err)
|
||||
}
|
||||
if fn.Status.Phase != slessv1alpha1.FunctionPhaseReady {
|
||||
fj.Status.Phase = slessv1alpha1.FunctionJobPhasePending
|
||||
fj.Status.Message = "waiting for function Ready (current: " + string(fn.Status.Phase) + ")"
|
||||
_ = r.Status().Update(ctx, fj)
|
||||
// RequeueAfter: опрашиваем каждые 15с пока Function не станет Ready.
|
||||
return ctrl.Result{RequeueAfter: 15 * time.Second}, nil
|
||||
// Фаза Building — ждём завершения kaniko Job
|
||||
if fj.Status.Phase == slessv1alpha1.FunctionJobPhaseBuilding {
|
||||
return r.checkJobBuild(ctx, fj)
|
||||
}
|
||||
|
||||
// Нет ImageRef — нужно собрать образ сначала.
|
||||
// Если аннотация build-job уже есть но phase не Building (рестарт контроллера),
|
||||
// восстанавливаем фазу Building.
|
||||
if fj.Status.ImageRef == "" {
|
||||
if fj.Annotations["sless.kube5s.ru/build-job"] != "" {
|
||||
fj.Status.Phase = slessv1alpha1.FunctionJobPhaseBuilding
|
||||
fj.Status.Message = "resuming build: " + fj.Annotations["sless.kube5s.ru/build-job"]
|
||||
_ = r.Status().Update(ctx, fj)
|
||||
return ctrl.Result{RequeueAfter: 10 * time.Second}, nil
|
||||
}
|
||||
return r.startJobBuild(ctx, fj)
|
||||
}
|
||||
|
||||
// Образ собран — создаём/синхронизируем k8s Job
|
||||
deployNS := "sless-fn-" + fj.Namespace
|
||||
jobName := fmt.Sprintf("job-%s-%s", fj.Name, fj.CreationTimestamp.Format("20060102150405"))
|
||||
|
||||
// Если Job уже создан — проверяем его статус
|
||||
existingJob := &batchv1.Job{}
|
||||
if err := r.Get(ctx, client.ObjectKey{Name: jobName, Namespace: deployNS}, existingJob); err == nil {
|
||||
return r.syncJobStatus(ctx, fj, existingJob)
|
||||
@@ -106,14 +106,105 @@ func (r *FunctionJobReconciler) Reconcile(ctx context.Context, req ctrl.Request)
|
||||
return ctrl.Result{}, fmt.Errorf("get job: %w", err)
|
||||
}
|
||||
|
||||
// Создаём k8s Job
|
||||
// Используем образ функции напрямую, переопределяем CMD чтобы запустить runner
|
||||
// вместо server.py/server.js — runner выполняет handle(event) один раз и выходит
|
||||
return r.createRunJob(ctx, fj, deployNS, jobName)
|
||||
}
|
||||
|
||||
// startJobBuild запускает kaniko сборку образа и переводит FunctionJob в фазу Building.
|
||||
func (r *FunctionJobReconciler) startJobBuild(ctx context.Context, fj *slessv1alpha1.FunctionJob) (ctrl.Result, error) {
|
||||
logger := log.FromContext(ctx)
|
||||
buildJobName, err := r.Builder.Build(ctx, r.OperatorNamespace, fj.Name, fj.Spec.S3Key)
|
||||
if err != nil {
|
||||
fj.Status.Phase = slessv1alpha1.FunctionJobPhaseFailed
|
||||
fj.Status.Message = "failed to start build: " + err.Error()
|
||||
_ = r.Status().Update(ctx, fj)
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
// Аннотация build-job guard против повторного запуска при параллельных reconcile.
|
||||
// Как только аннотация выставлена, следующий reconcile войдёт в checkJobBuild.
|
||||
if fj.Annotations == nil {
|
||||
fj.Annotations = map[string]string{}
|
||||
}
|
||||
fj.Annotations["sless.kube5s.ru/build-job"] = buildJobName
|
||||
if err := r.Update(ctx, fj); err != nil {
|
||||
return ctrl.Result{}, fmt.Errorf("update build annotation: %w", err)
|
||||
}
|
||||
|
||||
fj.Status.Phase = slessv1alpha1.FunctionJobPhaseBuilding
|
||||
fj.Status.Message = "building image: " + buildJobName
|
||||
if err := r.Status().Update(ctx, fj); err != nil {
|
||||
return ctrl.Result{}, fmt.Errorf("update status to building: %w", err)
|
||||
}
|
||||
|
||||
logger.Info("started kaniko build for functionjob", "build-job", buildJobName, "functionjob", fj.Name)
|
||||
return ctrl.Result{RequeueAfter: 10 * time.Second}, nil
|
||||
}
|
||||
|
||||
// checkJobBuild опрашивает статус kaniko Job.
|
||||
// При успехе: сохраняет ImageRef в status, очищает Build Job, requeue → createRunJob.
|
||||
// При ошибке: переводит FunctionJob в Failed с логами kaniko.
|
||||
func (r *FunctionJobReconciler) checkJobBuild(ctx context.Context, fj *slessv1alpha1.FunctionJob) (ctrl.Result, error) {
|
||||
logger := log.FromContext(ctx)
|
||||
buildJobName := fj.Annotations["sless.kube5s.ru/build-job"]
|
||||
if buildJobName == "" {
|
||||
fj.Status.Phase = slessv1alpha1.FunctionJobPhasePending
|
||||
_ = r.Status().Update(ctx, fj)
|
||||
return ctrl.Result{Requeue: true}, nil
|
||||
}
|
||||
|
||||
status, err := r.Builder.JobStatus(ctx, buildJobName)
|
||||
if err != nil {
|
||||
return ctrl.Result{}, fmt.Errorf("check build job: %w", err)
|
||||
}
|
||||
|
||||
switch status {
|
||||
case "running":
|
||||
return ctrl.Result{RequeueAfter: 10 * time.Second}, nil
|
||||
|
||||
case "succeeded":
|
||||
// Builder вычисляет imageRef детерминировано по namespace+name+s3Key
|
||||
imageRef := r.Builder.ImageRef(r.OperatorNamespace, fj.Name, fj.Spec.S3Key)
|
||||
fj.Status.ImageRef = imageRef
|
||||
// Сбрасываем Phase чтобы следующий reconcile пошёл в createRunJob
|
||||
fj.Status.Phase = slessv1alpha1.FunctionJobPhasePending
|
||||
fj.Status.Message = ""
|
||||
if err := r.Status().Update(ctx, fj); err != nil {
|
||||
return ctrl.Result{}, fmt.Errorf("update imageref in status: %w", err)
|
||||
}
|
||||
_ = r.Builder.Cleanup(ctx, buildJobName)
|
||||
logger.Info("build succeeded, queuing run job", "image", imageRef, "functionjob", fj.Name)
|
||||
return ctrl.Result{Requeue: true}, nil
|
||||
|
||||
case "failed":
|
||||
// Ищем поды kaniko по лейблу job-name в namespace оператора
|
||||
logs := getJobPodOutput(ctx, r.KubeClient, r.OperatorNamespace, "job-name="+buildJobName)
|
||||
msg := "build job failed"
|
||||
if logs != "" {
|
||||
msg = "build job failed:\n" + logs
|
||||
}
|
||||
fj.Status.Phase = slessv1alpha1.FunctionJobPhaseFailed
|
||||
fj.Status.Message = msg
|
||||
_ = r.Status().Update(ctx, fj)
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
return ctrl.Result{RequeueAfter: 10 * time.Second}, nil
|
||||
}
|
||||
|
||||
// createRunJob создаёт k8s Job для выполнения функции.
|
||||
// Использует собранный образ из fj.Status.ImageRef.
|
||||
func (r *FunctionJobReconciler) createRunJob(ctx context.Context, fj *slessv1alpha1.FunctionJob, deployNS, jobName string) (ctrl.Result, error) {
|
||||
logger := log.FromContext(ctx)
|
||||
eventJSON := fj.Spec.EventJSON
|
||||
if eventJSON == "" {
|
||||
eventJSON = "{}"
|
||||
}
|
||||
|
||||
memMB := fj.Spec.MemoryMB
|
||||
if memMB <= 0 {
|
||||
memMB = 128
|
||||
}
|
||||
|
||||
// runner запускается через env var SLESS_EVENT — безопаснее чем передавать в args
|
||||
// (args видны в ps aux, env vars — нет)
|
||||
ttl := int32(600) // автоудаление Job через 10 мин после завершения
|
||||
@@ -125,7 +216,6 @@ func (r *FunctionJobReconciler) Reconcile(ctx context.Context, req ctrl.Request)
|
||||
Labels: map[string]string{
|
||||
"managed-by": "sless",
|
||||
"functionjob": fj.Name,
|
||||
"function": fn.Name,
|
||||
},
|
||||
},
|
||||
Spec: batchv1.JobSpec{
|
||||
@@ -138,31 +228,25 @@ func (r *FunctionJobReconciler) Reconcile(ctx context.Context, req ctrl.Request)
|
||||
Labels: map[string]string{
|
||||
"managed-by": "sless",
|
||||
"functionjob": fj.Name,
|
||||
"function": fn.Name,
|
||||
},
|
||||
},
|
||||
Spec: corev1.PodSpec{
|
||||
RestartPolicy: corev1.RestartPolicyNever,
|
||||
// Используем тот же образ что и Deployment функции
|
||||
// runner.py/runner.js переопределяет CMD сервера
|
||||
InitContainers: nil,
|
||||
Containers: []corev1.Container{
|
||||
{
|
||||
Name: "runner",
|
||||
Image: fn.Status.ImageRef,
|
||||
// Переопределяем точку входа: запускаем runner вместо server
|
||||
// runner читает SLESS_EVENT и вызывает handle(event) один раз
|
||||
Command: runtimeRunnerCommand(fn.Spec.Runtime),
|
||||
Name: "runner",
|
||||
Image: fj.Status.ImageRef,
|
||||
Command: runtimeRunnerCommand(fj.Spec.Runtime),
|
||||
Env: append(
|
||||
append(fnEnvVars(fn), corev1.EnvVar{
|
||||
append(fjEnvVars(fj), corev1.EnvVar{
|
||||
Name: "SLESS_EVENT",
|
||||
Value: eventJSON,
|
||||
}),
|
||||
goJobModeEnv(fn.Spec.Runtime)...,
|
||||
goJobModeEnv(fj.Spec.Runtime)...,
|
||||
),
|
||||
Resources: corev1.ResourceRequirements{
|
||||
Limits: corev1.ResourceList{
|
||||
corev1.ResourceMemory: resource.MustParse(fmt.Sprintf("%dMi", fn.Spec.MemoryMB)),
|
||||
corev1.ResourceMemory: resource.MustParse(fmt.Sprintf("%dMi", memMB)),
|
||||
},
|
||||
},
|
||||
},
|
||||
@@ -188,10 +272,11 @@ func (r *FunctionJobReconciler) Reconcile(ctx context.Context, req ctrl.Request)
|
||||
return ctrl.Result{}, fmt.Errorf("update functionjob status: %w", err)
|
||||
}
|
||||
|
||||
logger.Info("created job for functionjob", "job", jobName, "functionjob", fj.Name)
|
||||
logger.Info("created run job for functionjob", "job", jobName, "functionjob", fj.Name)
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
|
||||
// syncJobStatus читает статус k8s Job и обновляет FunctionJob.Status.
|
||||
// Если Job ещё выполняется — запрашивает повторный reconcile через 5 сек
|
||||
// (Owns не работает кросс-неймспейсно, поэтому используем polling).
|
||||
@@ -282,13 +367,13 @@ print(json.dumps(result))
|
||||
}
|
||||
}
|
||||
|
||||
// fnEnvVars преобразует env vars из FunctionSpec в k8s EnvVar slice.
|
||||
// Включает SLESS_ENTRYPOINT чтобы runner.py/runner.js знал какую функцию вызывать.
|
||||
func fnEnvVars(fn *slessv1alpha1.Function) []corev1.EnvVar {
|
||||
// fjEnvVars формирует k8s EnvVar из полей FunctionJobSpec.
|
||||
// SLESS_ENTRYPOINT сообщает runner'у какую функцию вызывать.
|
||||
func fjEnvVars(fj *slessv1alpha1.FunctionJob) []corev1.EnvVar {
|
||||
result := []corev1.EnvVar{
|
||||
{Name: "SLESS_ENTRYPOINT", Value: fn.Spec.Entrypoint},
|
||||
{Name: "SLESS_ENTRYPOINT", Value: fj.Spec.Entrypoint},
|
||||
}
|
||||
for k, v := range fn.Spec.Env {
|
||||
for k, v := range fj.Spec.Env {
|
||||
result = append(result, corev1.EnvVar{Name: k, Value: v})
|
||||
}
|
||||
return result
|
||||
|
||||
@@ -2,6 +2,48 @@
|
||||
|
||||
---
|
||||
|
||||
## 2026-03-20 — Merge: убрать sless_function как обязательный prerequisite для sless_job
|
||||
|
||||
### Контекст
|
||||
|
||||
`sless_job` ранее требовал `FunctionRef` — имя существующего `sless_function` из которого брался `ImageRef`.
|
||||
Это создавало два отдельных ресурса для одной задачи (запустить код один раз):
|
||||
|
||||
```hcl
|
||||
resource "sless_function" "f" { ... } # build
|
||||
resource "sless_job" "j" { function = sless_function.f.name ... } # run
|
||||
```
|
||||
|
||||
### Решение
|
||||
|
||||
Сделать `FunctionJobSpec` самодостаточным: встроить `Runtime/Entrypoint/Env/S3Key` и запускать
|
||||
kaniko сборку непосредственно из FunctionJob-контроллера (новая фаза `Building`).
|
||||
|
||||
```hcl
|
||||
resource "sless_job" "j" {
|
||||
runtime = "python3.11"
|
||||
source_dir = "./code/fn"
|
||||
...
|
||||
}
|
||||
```
|
||||
|
||||
### Почему НЕ удаляем sless_function
|
||||
|
||||
`sless_function` нужен для `sless_trigger` (type=cron/http) — они ссылаются на функцию.
|
||||
Для триггеров образ должен жить вечно (не удаляться после запуска), и за ним следит Function CRD.
|
||||
`sless_job` же — разовый запуск; после завершения Job удаляется, образ остаётся в registry.
|
||||
|
||||
### Изменения в State Machine FunctionJobReconciler
|
||||
|
||||
```
|
||||
Было: Pending → (ждать Function.Ready) → Running → Succeeded/Failed
|
||||
Стало: Pending → Building (kaniko) → Pending + ImageRef → Running → Succeeded/Failed
|
||||
```
|
||||
|
||||
Фаза `Building` охраняется аннотацией `sless.kube5s.ru/build-job` — идемпотентна при рестарте.
|
||||
|
||||
---
|
||||
|
||||
## 2026-03-19 — Go runtime v0.1.1: внешние зависимости через go.mod/go.sum
|
||||
|
||||
### Контекст
|
||||
|
||||
+29
-1
@@ -1,9 +1,37 @@
|
||||
# Прогресс разработки
|
||||
|
||||
Последнее обновление: 2026-03-20 21:30
|
||||
Последнее обновление: 2026-03-20 (merge: sless_function + sless_job → self-contained sless_job)
|
||||
|
||||
---
|
||||
|
||||
## 2026-03-20 — Merge: sless_function + sless_job → единый self-contained sless_job
|
||||
|
||||
### Цель
|
||||
Убрать обязательную зависимость `sless_job` от `sless_function`. Раньше для запуска одноразового
|
||||
джоба нужно было сначала создать `sless_function` (CRD + kaniko build), потом `sless_job`.
|
||||
Теперь `sless_job` самодостаточен: содержит runtime/entrypoint/env/source_dir и сам запускает сборку.
|
||||
|
||||
### Изменённые файлы
|
||||
|
||||
| Файл | Что сделано |
|
||||
|------|-------------|
|
||||
| `api/v1alpha1/job_types.go` | Убран `FunctionRef`. Добавлены: `Runtime/Entrypoint/S3Bucket/S3Key/MemoryMB/TimeoutSec/Env map[string]string`. Фаза `Building`. `ImageRef` в status. |
|
||||
| `api/v1alpha1/zz_generated.deepcopy.go` | DeepCopyInto для FunctionJobSpec: proper deep copy map[string]string Env |
|
||||
| `controllers/functionjob_controller.go` | Полная переработка: убрана зависимость от Function CRD. Новые поля Builder+OperatorNamespace. Новая фаза Building (kaniko). Методы startJobBuild/checkJobBuild. |
|
||||
| `main.go` | Передача `Builder: bldr` и `OperatorNamespace: "sless"` в FunctionJobReconciler |
|
||||
| `internal/api/handler/jobs.go` | jobRequest/jobResponse без FunctionRef, с Runtime/Entrypoint/Env/S3Key. Новый handler UploadJobCode. |
|
||||
| `internal/api/router.go` | Добавлен маршрут `/namespaces/{ns}/jobs/{name}/upload` |
|
||||
| `terraform/provider/internal/client/client.go` | JobRequest/JobResponse без FunctionRef, с Runtime/Entrypoint/Env/ImageRef. UploadJobCode метод. Рефакторинг UploadCodeReader → uploadCodeToURL. |
|
||||
| `terraform/provider/internal/resources/job_resource.go` | JobModel без Function, с Runtime/Entrypoint/EnvVars/SourceDir/CodeHash/ImageRef. ModifyPlan. Create с upload. Дефолт wait_timeout_sec=900. |
|
||||
| `examples/POSTGRES/functions.tf` | Убран `sless_function.postgres_sql_runner_create_table`. `sless_job.postgres_table_init_job` теперь самодостаточен: inline source_dir/runtime/entrypoint/env_vars. |
|
||||
|
||||
### Статус
|
||||
- ✅ Go код скомпилировался (controllers, main, internal/api)
|
||||
- ✅ Terraform provider: нет ошибок компилятора
|
||||
- ✅ functions.tf: раскомментирован и обновлён
|
||||
- ⏳ Требует: кросс-компиляция provider → deploy в кластер → тест apply
|
||||
|
||||
|
||||
## 2026-03-20 — Восстановление PostgreSQL и отладка провайдера nubes
|
||||
|
||||
### Что произошло
|
||||
|
||||
@@ -1,22 +1,18 @@
|
||||
// 2026-03-20 — удалены stress_* функции (oneshot/Job). Архив: git history.
|
||||
// 2026-03-20 — выделено из resources.tf: sless функции, сервисы и джобы.
|
||||
// 2026-03-19 — миграция: sless_function (oneshot/Job) + sless_service (long-running Deployment).
|
||||
// sless_trigger(type=http) удалены — HTTP URL теперь автоматически в sless_service.
|
||||
// postgres_sql_runner_create_table остаётся sless_function (oneshot/Job).
|
||||
// pg-info, pg-table-reader, pg-table-writer → sless_service.
|
||||
// 2026-03-20 — все ресурсы закомментированы перед удалением инстанса PostgreSQL.
|
||||
// 2026-03-20 (merge: sless_function + старый sless_job объединены в один self-contained sless_job)
|
||||
// Теперь sless_job несёт в себе runtime/entrypoint/source_dir — не нужен отдельный sless_function.
|
||||
// WaitJobDone таймаут 900s покрывает kaniko сборку (~5 мин) + выполнение SQL (~несколько сек).
|
||||
|
||||
/*
|
||||
|
||||
# Служебная функция выполняет SQL-операторы из event_json.
|
||||
# Credentials берутся из locals (vault_secrets) — без хардкода.
|
||||
# Для сверки хардкод остаётся в terraform.tfvars.
|
||||
resource "sless_function" "postgres_sql_runner_create_table" {
|
||||
name = "pg-create-table-runner"
|
||||
runtime = "python3.11"
|
||||
entrypoint = "sql_runner.run_sql"
|
||||
memory_mb = 128
|
||||
timeout_sec = 30
|
||||
# Одноразовый запуск: собирает образ через kaniko, выполняет SQL, завершается.
|
||||
# Заменяет sless_function.postgres_sql_runner_create_table + sless_job.postgres_table_init_job.
|
||||
resource "sless_job" "postgres_table_init_job" {
|
||||
name = "pg-create-table-job-main-v13"
|
||||
runtime = "python3.11"
|
||||
entrypoint = "sql_runner.run_sql"
|
||||
memory_mb = 128
|
||||
timeout_sec = 30
|
||||
source_dir = "${path.module}/code/sql-runner"
|
||||
wait_timeout_sec = 900
|
||||
run_id = 13
|
||||
|
||||
env_vars = {
|
||||
PGHOST = local.pg_host
|
||||
@@ -25,20 +21,8 @@ resource "sless_function" "postgres_sql_runner_create_table" {
|
||||
PGUSER = local.pg_username
|
||||
PGPASSWORD = local.pg_password
|
||||
PGSSLMODE = "require"
|
||||
# Для сверки (должно совпадать с vault):
|
||||
# PGUSER = var.pg_user
|
||||
# PGPASSWORD = var.pg_password
|
||||
}
|
||||
|
||||
source_dir = "${path.module}/code/sql-runner"
|
||||
}
|
||||
|
||||
resource "sless_job" "postgres_table_init_job" {
|
||||
name = "pg-create-table-job-main-v13"
|
||||
function = sless_function.postgres_sql_runner_create_table.name
|
||||
wait_timeout_sec = 180
|
||||
run_id = 13
|
||||
|
||||
event_json = jsonencode({
|
||||
statements = [
|
||||
"CREATE TABLE IF NOT EXISTS terraform_demo_table (id serial PRIMARY KEY, title text NOT NULL, created_at timestamp DEFAULT now())"
|
||||
@@ -49,8 +33,6 @@ resource "sless_job" "postgres_table_init_job" {
|
||||
}
|
||||
|
||||
# Long-running сервис на NodeJS: возвращает версию PG-сервера и счётчик строк в таблице.
|
||||
# Единственная функция примера на nodejs20 — проверка что JS runtime работает.
|
||||
# URL автоматически: https://sless.kube5s.ru/fn/<namespace>/pg-info
|
||||
resource "sless_service" "pg_info" {
|
||||
name = "pg-info"
|
||||
runtime = "nodejs20"
|
||||
@@ -72,10 +54,6 @@ resource "sless_service" "pg_info" {
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
|
||||
# Long-running сервисы чтения и записи строк terraform_demo_table — в одном файле table_rw.py.
|
||||
# list_rows (GET) — читает все строки; add_row (POST {title}) — вставляет строку.
|
||||
# URL автоматически: https://sless.kube5s.ru/fn/<namespace>/pg-table-reader
|
||||
# https://sless.kube5s.ru/fn/<namespace>/pg-table-writer
|
||||
resource "sless_service" "postgres_table_reader" {
|
||||
name = "pg-table-reader"
|
||||
runtime = "python3.11"
|
||||
@@ -125,5 +103,3 @@ resource "sless_service" "postgres_table_writer" {
|
||||
output "table_writer_url" {
|
||||
value = sless_service.postgres_table_writer.url
|
||||
}
|
||||
|
||||
*/
|
||||
|
||||
+110
-19
@@ -1,4 +1,4 @@
|
||||
// Изменено: 2026-03-08
|
||||
// Изменено: 2026-03-20 (merge: FunctionJob теперь самодостаточен — убран FunctionRef, добавлены Runtime/Entrypoint/Env)
|
||||
// jobs.go — CRUD handlers для FunctionJob CRD.
|
||||
// Создаёт/читает/удаляет k8s FunctionJob ресурсы.
|
||||
// Namespace берётся из URL: /v1/namespaces/{namespace}/jobs/{name}
|
||||
@@ -7,20 +7,29 @@ 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"
|
||||
)
|
||||
|
||||
// jobRequest — тело POST /v1/namespaces/{ns}/jobs
|
||||
type jobRequest struct {
|
||||
Name string `json:"name"`
|
||||
FunctionRef string `json:"function"`
|
||||
EventJSON string `json:"event_json,omitempty"`
|
||||
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"`
|
||||
}
|
||||
@@ -29,10 +38,12 @@ type jobRequest struct {
|
||||
type jobResponse struct {
|
||||
Name string `json:"name"`
|
||||
Namespace string `json:"namespace"`
|
||||
FunctionRef string `json:"function"`
|
||||
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"`
|
||||
@@ -42,14 +53,16 @@ type jobResponse struct {
|
||||
// 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,
|
||||
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")
|
||||
@@ -61,7 +74,7 @@ func jobToResponse(j *slessv1alpha1.FunctionJob) jobResponse {
|
||||
}
|
||||
|
||||
// CreateJob — POST /v1/namespaces/{namespace}/jobs
|
||||
// Создаёт FunctionJob CR. Оператор запустит k8s Job асинхронно.
|
||||
// Создаёт FunctionJob CR. Оператор запустит kaniko сборку и затем k8s Job асинхронно.
|
||||
func (h *Handler) CreateJob(w http.ResponseWriter, r *http.Request) {
|
||||
ns := namespace(r)
|
||||
|
||||
@@ -74,8 +87,12 @@ func (h *Handler) CreateJob(w http.ResponseWriter, r *http.Request) {
|
||||
writeJSON(w, http.StatusBadRequest, errResp("name is required"))
|
||||
return
|
||||
}
|
||||
if req.FunctionRef == "" {
|
||||
writeJSON(w, http.StatusBadRequest, errResp("function is required"))
|
||||
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 == "" {
|
||||
@@ -88,9 +105,15 @@ func (h *Handler) CreateJob(w http.ResponseWriter, r *http.Request) {
|
||||
Namespace: ns,
|
||||
},
|
||||
Spec: slessv1alpha1.FunctionJobSpec{
|
||||
FunctionRef: req.FunctionRef,
|
||||
EventJSON: req.EventJSON,
|
||||
RunID: req.RunID,
|
||||
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,
|
||||
},
|
||||
}
|
||||
|
||||
@@ -152,3 +175,71 @@ func (h *Handler) DeleteJob(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
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
|
||||
fj := &slessv1alpha1.FunctionJob{}
|
||||
if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, fj); err != nil {
|
||||
if errors.IsNotFound(err) {
|
||||
writeJSON(w, http.StatusNotFound, errResp("job not found"))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusInternalServerError, errResp(err.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
|
||||
}
|
||||
|
||||
// Версия на основе timestamp — каждый upload → новый уникальный ключ в S3
|
||||
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
|
||||
}
|
||||
|
||||
// Обновляем 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",
|
||||
})
|
||||
}
|
||||
|
||||
@@ -67,6 +67,7 @@ func NewRouter(h *handler.Handler, log *slog.Logger) http.Handler {
|
||||
v1.HandleFunc("/namespaces/{namespace}/jobs", h.CreateJob).Methods(http.MethodPost)
|
||||
v1.HandleFunc("/namespaces/{namespace}/jobs/{name}", h.GetJob).Methods(http.MethodGet)
|
||||
v1.HandleFunc("/namespaces/{namespace}/jobs/{name}", h.DeleteJob).Methods(http.MethodDelete)
|
||||
v1.HandleFunc("/namespaces/{namespace}/jobs/{name}/upload", h.UploadJobCode).Methods(http.MethodPost)
|
||||
|
||||
// Цепочка middleware: logging → (auth только для /v1/) → router
|
||||
// /fn/ — без auth, /v1/ — с auth.
|
||||
|
||||
@@ -178,10 +178,12 @@ func main() {
|
||||
os.Exit(1)
|
||||
}
|
||||
if err = (&controllers.FunctionJobReconciler{
|
||||
Client: mgr.GetClient(),
|
||||
Scheme: mgr.GetScheme(),
|
||||
RegistrySecret: cfg.RegistrySecret,
|
||||
KubeClient: kubernetes.NewForConfigOrDie(mgr.GetConfig()),
|
||||
Client: mgr.GetClient(),
|
||||
Scheme: mgr.GetScheme(),
|
||||
RegistrySecret: cfg.RegistrySecret,
|
||||
KubeClient: kubernetes.NewForConfigOrDie(mgr.GetConfig()),
|
||||
Builder: bldr,
|
||||
OperatorNamespace: "sless",
|
||||
}).SetupWithManager(mgr); err != nil {
|
||||
log.Error("unable to create controller", "controller", "FunctionJob", "err", err)
|
||||
os.Exit(1)
|
||||
|
||||
@@ -298,6 +298,17 @@ func (c *Client) UploadCode(ctx context.Context, ns, name, zipPath string) error
|
||||
// UploadCodeReader — загружает код из произвольного io.Reader (например in-memory zip).
|
||||
// filename используется только как имя файла в multipart-форме.
|
||||
func (c *Client) UploadCodeReader(ctx context.Context, ns, name, filename string, r io.Reader) error {
|
||||
return c.uploadCodeToURL(ctx, fmt.Sprintf("%s/v1/namespaces/%s/functions/%s/upload", c.endpoint, ns, name), filename, r)
|
||||
}
|
||||
|
||||
// UploadJobCode — POST /v1/namespaces/{ns}/jobs/{name}/upload
|
||||
// Аналогично UploadCodeReader но для FunctionJob — запускает kaniko через FunctionJob CRD.
|
||||
func (c *Client) UploadJobCode(ctx context.Context, ns, name, filename string, r io.Reader) error {
|
||||
return c.uploadCodeToURL(ctx, fmt.Sprintf("%s/v1/namespaces/%s/jobs/%s/upload", c.endpoint, ns, name), filename, r)
|
||||
}
|
||||
|
||||
// uploadCodeToURL — внутренний хелпер: пакует io.Reader в multipart и POST-ит по указанному URL.
|
||||
func (c *Client) uploadCodeToURL(ctx context.Context, uploadURL, filename string, r io.Reader) error {
|
||||
var buf bytes.Buffer
|
||||
mw := multipart.NewWriter(&buf)
|
||||
fw, err := mw.CreateFormFile("code", filename)
|
||||
@@ -309,8 +320,7 @@ func (c *Client) UploadCodeReader(ctx context.Context, ns, name, filename string
|
||||
}
|
||||
mw.Close()
|
||||
|
||||
url := fmt.Sprintf("%s/v1/namespaces/%s/functions/%s/upload", c.endpoint, ns, name)
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, &buf)
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, uploadURL, &buf)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -428,10 +438,17 @@ func (c *Client) UpdateTrigger(ctx context.Context, ns, name string, req Trigger
|
||||
// --- Job CRUD ---
|
||||
|
||||
// JobRequest — тело POST /v1/namespaces/{ns}/jobs
|
||||
// Изменено: 2026-03-20 (merge: убран FunctionRef, добавлены Runtime/Entrypoint/Env)
|
||||
type JobRequest struct {
|
||||
Name string `json:"name"`
|
||||
FunctionRef string `json:"function"`
|
||||
EventJSON string `json:"event_json,omitempty"`
|
||||
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"`
|
||||
}
|
||||
@@ -440,10 +457,12 @@ type JobRequest struct {
|
||||
type JobResponse struct {
|
||||
Name string `json:"name"`
|
||||
Namespace string `json:"namespace"`
|
||||
FunctionRef string `json:"function"`
|
||||
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"`
|
||||
JobName string `json:"job_name"`
|
||||
StartTime string `json:"start_time"`
|
||||
CompletionTime string `json:"completion_time"`
|
||||
|
||||
@@ -1,21 +1,23 @@
|
||||
// 2026-03-17 12:20
|
||||
// 2026-03-20 (merge: убран FunctionRef, добавлены Runtime/Entrypoint/SourceDir — sless_job теперь самодостаточен)
|
||||
// job_resource.go — Terraform ресурс sless_job.
|
||||
//
|
||||
// Lifecycle:
|
||||
//
|
||||
// Create: POST /v1/namespaces/{ns}/jobs → WaitJobDone (10 мин)
|
||||
// Блокирует terraform apply до завершения джоба (Succeeded/Failed).
|
||||
// Если Failed — terraform apply падает с ошибкой.
|
||||
// Если run_id=0 — джоб создаётся в k8s, но k8s Job не запускается.
|
||||
// Create: POST /v1/namespaces/{ns}/jobs
|
||||
// Если source_dir задан → zipDir → UploadJobCode → контроллер запускает kaniko
|
||||
// WaitJobDone покрывает: Building (kaniko) + Running (функция) + Succeeded/Failed
|
||||
// Если run_id=0 — джоб создаётся в k8s, но не запускается (phase=Skipped).
|
||||
// Read: GET /v1/namespaces/{ns}/jobs/{name} → sync phase/timing в state
|
||||
// Delete: DELETE /v1/namespaces/{ns}/jobs/{name}
|
||||
// Семантически no-op (джоб уже выполнен), но убирает CR из кластера.
|
||||
//
|
||||
// run_id: значение 0 = создать без запуска. >0 = запустить/перезапустить джоб.
|
||||
// run_id имеет RequiresReplace: изменение значения (1→2→3) триггерирует повторный запуск.
|
||||
// run_id: 0 = создать без запуска; >0 = запустить.
|
||||
// RequiresReplace на run_id + code_hash: изменение кода или run_id всегда пересоздаёт джоб.
|
||||
// wait_timeout_sec должен покрывать оба этапа: kaniko (~5 мин) + выполнение функции.
|
||||
// По умолчанию 900 сек (15 мин).
|
||||
package resources
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
@@ -34,11 +36,13 @@ import (
|
||||
"github.com/hashicorp/terraform-plugin-framework/types"
|
||||
)
|
||||
|
||||
// defaultWaitTimeoutSec — дефолтный таймаут ожидания завершения джоба (600 сек = 10 мин).
|
||||
// Пользователь может переопределить через wait_timeout_sec.
|
||||
const defaultWaitTimeoutSec = 600
|
||||
// defaultJobWaitTimeoutSec — дефолтный таймаут ожидания завершения sless_job.
|
||||
// Покрывает kaniko сборку (~5 мин) + выполнение функции (~остаток).
|
||||
// Переопределяется через wait_timeout_sec.
|
||||
const defaultJobWaitTimeoutSec = 900
|
||||
|
||||
var _ resource.Resource = &JobResource{}
|
||||
var _ resource.ResourceWithModifyPlan = &JobResource{}
|
||||
|
||||
type JobResource struct {
|
||||
client *client.Client
|
||||
@@ -50,15 +54,24 @@ func NewJobResource() resource.Resource {
|
||||
|
||||
// JobModel — модель состояния terraform для sless_job.
|
||||
type JobModel struct {
|
||||
Name types.String `tfsdk:"name"`
|
||||
Function types.String `tfsdk:"function"`
|
||||
Name types.String `tfsdk:"name"`
|
||||
Runtime types.String `tfsdk:"runtime"`
|
||||
Entrypoint types.String `tfsdk:"entrypoint"`
|
||||
MemoryMB types.Int64 `tfsdk:"memory_mb"`
|
||||
TimeoutSec types.Int64 `tfsdk:"timeout_sec"`
|
||||
EnvVars types.Map `tfsdk:"env_vars"`
|
||||
// source_dir — директория с исходниками. Провайдер сам упакует в zip и загрузит.
|
||||
SourceDir types.String `tfsdk:"source_dir"`
|
||||
// code_hash — SHA256 содержимого source_dir. Computed; при изменении RequiresReplace.
|
||||
CodeHash types.String `tfsdk:"code_hash"`
|
||||
EventJSON types.String `tfsdk:"event_json"`
|
||||
// RunID: 0 = не запускать (Skipped), 1+ = запустить/перезапустить.
|
||||
// RequiresReplace: изменение = пересоздание FunctionJob → новый запуск.
|
||||
RunID types.Int64 `tfsdk:"run_id"`
|
||||
// wait_timeout_sec — максимальное ожидание завершения джоба. Дефолт 600 сек.
|
||||
// wait_timeout_sec — максимальное ожидание (kaniko + выполнение). Дефолт 900 сек.
|
||||
WaitTimeoutSec types.Int64 `tfsdk:"wait_timeout_sec"`
|
||||
Phase types.String `tfsdk:"phase"`
|
||||
ImageRef types.String `tfsdk:"image_ref"`
|
||||
StartTime types.String `tfsdk:"start_time"`
|
||||
CompletionTime types.String `tfsdk:"completion_time"`
|
||||
Message types.String `tfsdk:"message"`
|
||||
@@ -70,19 +83,56 @@ func (r *JobResource) Metadata(_ context.Context, req resource.MetadataRequest,
|
||||
|
||||
func (r *JobResource) Schema(_ context.Context, _ resource.SchemaRequest, resp *resource.SchemaResponse) {
|
||||
resp.Schema = schema.Schema{
|
||||
MarkdownDescription: "Одноразовый запуск serverless функции. terraform apply блокируется до завершения джоба.",
|
||||
MarkdownDescription: "Одноразовый запуск serverless функции. terraform apply блокируется до завершения. Включает kaniko сборку образа.",
|
||||
Attributes: map[string]schema.Attribute{
|
||||
// Все input-поля immutable — джоб нельзя "изменить", только пересоздать.
|
||||
|
||||
"name": schema.StringAttribute{
|
||||
Required: true,
|
||||
PlanModifiers: []planmodifier.String{
|
||||
stringplanmodifier.RequiresReplace(),
|
||||
},
|
||||
},
|
||||
"function": schema.StringAttribute{
|
||||
"runtime": schema.StringAttribute{
|
||||
Required: true,
|
||||
MarkdownDescription: "Имя sless_function ресурса в том же namespace.",
|
||||
MarkdownDescription: "Среда выполнения: python3.11, nodejs20, go1.23.",
|
||||
PlanModifiers: []planmodifier.String{
|
||||
stringplanmodifier.RequiresReplace(),
|
||||
},
|
||||
},
|
||||
"entrypoint": schema.StringAttribute{
|
||||
Required: true,
|
||||
MarkdownDescription: "Точка входа: module.function (например sql_runner.run_sql).",
|
||||
PlanModifiers: []planmodifier.String{
|
||||
stringplanmodifier.RequiresReplace(),
|
||||
},
|
||||
},
|
||||
"memory_mb": schema.Int64Attribute{
|
||||
Optional: true,
|
||||
Computed: true,
|
||||
Default: int64default.StaticInt64(128),
|
||||
MarkdownDescription: "Лимит оперативной памяти в MB. По умолчанию 128.",
|
||||
},
|
||||
"timeout_sec": schema.Int64Attribute{
|
||||
Optional: true,
|
||||
Computed: true,
|
||||
Default: int64default.StaticInt64(30),
|
||||
MarkdownDescription: "Таймаут выполнения функции в секундах. По умолчанию 30.",
|
||||
},
|
||||
"env_vars": schema.MapAttribute{
|
||||
ElementType: types.StringType,
|
||||
Optional: true,
|
||||
MarkdownDescription: "Переменные окружения для функции.",
|
||||
},
|
||||
"source_dir": schema.StringAttribute{
|
||||
Optional: true,
|
||||
MarkdownDescription: "Директория с исходниками. Провайдер сам запакует zip и загрузит.",
|
||||
PlanModifiers: []planmodifier.String{
|
||||
stringplanmodifier.RequiresReplace(),
|
||||
},
|
||||
},
|
||||
// code_hash — вычисляется автоматически из source_dir (ModifyPlan).
|
||||
// RequiresReplace: если код изменился (hash != state) → пересоздать джоб.
|
||||
"code_hash": schema.StringAttribute{
|
||||
Computed: true,
|
||||
PlanModifiers: []planmodifier.String{
|
||||
stringplanmodifier.RequiresReplace(),
|
||||
},
|
||||
@@ -93,29 +143,36 @@ func (r *JobResource) Schema(_ context.Context, _ resource.SchemaRequest, resp *
|
||||
PlanModifiers: []planmodifier.String{
|
||||
stringplanmodifier.RequiresReplace(),
|
||||
},
|
||||
}, // run_id: 0 = создать без запуска (Skipped), >0 = запустить.
|
||||
},
|
||||
// run_id: 0 = создать без запуска (Skipped), >0 = запустить.
|
||||
// Изменение run_id (1→2→3...) триггерирует пересоздание = новый запуск.
|
||||
"run_id": schema.Int64Attribute{
|
||||
Optional: true,
|
||||
Computed: true,
|
||||
Default: int64default.StaticInt64(0),
|
||||
MarkdownDescription: "0 = не запускать; >0 = запустить. Увеличьте run_id для повторного запуска джоба.",
|
||||
MarkdownDescription: "0 = не запускать; >0 = запустить. Увеличьте run_id для повторного запуска.",
|
||||
PlanModifiers: []planmodifier.Int64{
|
||||
int64planmodifier.RequiresReplace(),
|
||||
},
|
||||
Validators: []validator.Int64{
|
||||
int64validator.AtLeast(0),
|
||||
},
|
||||
}, // wait_timeout_sec — сколько ждать завершения джоба. Увеличь если код долго работает (например миграция БД).
|
||||
},
|
||||
// wait_timeout_sec покрывает kaniko сборку + выполнение функции.
|
||||
"wait_timeout_sec": schema.Int64Attribute{
|
||||
Optional: true,
|
||||
Computed: true,
|
||||
MarkdownDescription: "Таймаут ожидания завершения джоба в секундах. По умолчанию 600 (10 мин).",
|
||||
Default: int64default.StaticInt64(defaultJobWaitTimeoutSec),
|
||||
MarkdownDescription: "Таймаут ожидания (kaniko + выполнение) в секундах. По умолчанию 900.",
|
||||
},
|
||||
// Computed — заполняются после завершения джоба
|
||||
// Computed — заполняются контроллером и финализируются после завершения
|
||||
"phase": schema.StringAttribute{
|
||||
Computed: true,
|
||||
MarkdownDescription: "Фаза выполнения: Pending, Running, Succeeded, Failed.",
|
||||
MarkdownDescription: "Фаза: Pending, Building, Running, Succeeded, Failed.",
|
||||
},
|
||||
"image_ref": schema.StringAttribute{
|
||||
Computed: true,
|
||||
MarkdownDescription: "Docker image собранный kaniko.",
|
||||
},
|
||||
"start_time": schema.StringAttribute{
|
||||
Computed: true,
|
||||
@@ -161,24 +218,32 @@ func (r *JobResource) Create(ctx context.Context, req resource.CreateRequest, re
|
||||
eventJSON = "{}"
|
||||
}
|
||||
|
||||
// run_id=0: создаём FunctionJob в k8s, но оператор не запустит k8s Job.
|
||||
// Пользователь может поменять run_id > 0 позже чтобы запустить.
|
||||
envMap, diags := mapToStringMap(ctx, plan.EnvVars)
|
||||
resp.Diagnostics.Append(diags...)
|
||||
if resp.Diagnostics.HasError() {
|
||||
return
|
||||
}
|
||||
|
||||
runID := plan.RunID.ValueInt64()
|
||||
createdJob, err := r.client.CreateJob(ctx, ns, client.JobRequest{
|
||||
Name: plan.Name.ValueString(),
|
||||
FunctionRef: plan.Function.ValueString(),
|
||||
EventJSON: eventJSON,
|
||||
RunID: runID,
|
||||
Name: plan.Name.ValueString(),
|
||||
Runtime: plan.Runtime.ValueString(),
|
||||
Entrypoint: plan.Entrypoint.ValueString(),
|
||||
MemoryMB: int32(plan.MemoryMB.ValueInt64()),
|
||||
TimeoutSec: int32(plan.TimeoutSec.ValueInt64()),
|
||||
Env: envMap,
|
||||
EventJSON: eventJSON,
|
||||
RunID: runID,
|
||||
})
|
||||
if err != nil {
|
||||
if errors.Is(err, client.ErrJobAlreadyExists) {
|
||||
existingJob, getErr := r.client.GetJob(ctx, ns, plan.Name.ValueString())
|
||||
if getErr != nil {
|
||||
resp.Diagnostics.AddError("create job", fmt.Sprintf("job already exists and get existing failed: %s", getErr.Error()))
|
||||
resp.Diagnostics.AddError("create job", fmt.Sprintf("job already exists and get failed: %s", getErr.Error()))
|
||||
return
|
||||
}
|
||||
if existingJob == nil {
|
||||
resp.Diagnostics.AddError("create job", "job already exists but cannot be read right after conflict")
|
||||
resp.Diagnostics.AddError("create job", "job already exists but cannot be read after conflict")
|
||||
return
|
||||
}
|
||||
createdJob = existingJob
|
||||
@@ -188,31 +253,30 @@ func (r *JobResource) Create(ctx context.Context, req resource.CreateRequest, re
|
||||
}
|
||||
}
|
||||
|
||||
// Если RunID=0 — не ждём завершения, пишем state сразу
|
||||
if runID == 0 {
|
||||
// state: Namespace/Name/Function из plan, RunID=0, Phase=Skipped, остальное empty
|
||||
waitTimeoutSec := plan.WaitTimeoutSec
|
||||
if waitTimeoutSec.IsNull() || waitTimeoutSec.IsUnknown() || waitTimeoutSec.ValueInt64() <= 0 {
|
||||
waitTimeoutSec = types.Int64Value(defaultWaitTimeoutSec)
|
||||
// Загружаем код если задан source_dir — контроллер начнёт kaniko сборку после upload
|
||||
if !plan.SourceDir.IsNull() && plan.SourceDir.ValueString() != "" {
|
||||
zipData, hash, err := zipDir(plan.SourceDir.ValueString())
|
||||
if err != nil {
|
||||
resp.Diagnostics.AddError("zip source_dir", err.Error())
|
||||
return
|
||||
}
|
||||
resp.Diagnostics.Append(resp.State.Set(ctx, JobModel{
|
||||
Name: types.StringValue(plan.Name.ValueString()),
|
||||
Function: types.StringValue(plan.Function.ValueString()),
|
||||
EventJSON: plan.EventJSON,
|
||||
RunID: types.Int64Value(0),
|
||||
WaitTimeoutSec: waitTimeoutSec,
|
||||
Phase: types.StringValue(createdJob.Phase),
|
||||
StartTime: types.StringValue(createdJob.StartTime),
|
||||
CompletionTime: types.StringValue(createdJob.CompletionTime),
|
||||
Message: types.StringValue(createdJob.Message),
|
||||
})...)
|
||||
if err := r.client.UploadJobCode(ctx, ns, plan.Name.ValueString(), "function.zip", bytes.NewReader(zipData)); err != nil {
|
||||
resp.Diagnostics.AddError("upload job code", err.Error())
|
||||
return
|
||||
}
|
||||
plan.CodeHash = types.StringValue(hash)
|
||||
}
|
||||
|
||||
// run_id=0: не ждём, пишем state сразу (job в состоянии Skipped)
|
||||
if runID == 0 {
|
||||
resp.Diagnostics.Append(resp.State.Set(ctx, jobToModel(plan, createdJob))...)
|
||||
return
|
||||
}
|
||||
|
||||
// Блокируем apply до завершения джоба (Succeeded или Failed)
|
||||
// Блокируем apply до завершения джоба (охватывает Building + Running → Succeeded/Failed)
|
||||
waitSec := plan.WaitTimeoutSec.ValueInt64()
|
||||
if waitSec <= 0 {
|
||||
waitSec = defaultWaitTimeoutSec
|
||||
waitSec = defaultJobWaitTimeoutSec
|
||||
}
|
||||
j, err := r.client.WaitJobDone(ctx, ns, plan.Name.ValueString(), time.Duration(waitSec)*time.Second)
|
||||
if err != nil {
|
||||
@@ -236,7 +300,6 @@ func (r *JobResource) Read(ctx context.Context, req resource.ReadRequest, resp *
|
||||
return
|
||||
}
|
||||
if j == nil {
|
||||
// Джоб удалён вне terraform — убираем из state
|
||||
resp.State.RemoveResource(ctx)
|
||||
return
|
||||
}
|
||||
@@ -244,10 +307,10 @@ func (r *JobResource) Read(ctx context.Context, req resource.ReadRequest, resp *
|
||||
resp.Diagnostics.Append(resp.State.Set(ctx, jobToModel(state, j))...)
|
||||
}
|
||||
|
||||
// Update не реализован — все поля поддерживают только RequiresReplace.
|
||||
// Update не реализован — все input-поля имеют RequiresReplace.
|
||||
// terraform-plugin-framework никогда не вызовет Update для этого ресурса.
|
||||
func (r *JobResource) Update(_ context.Context, _ resource.UpdateRequest, resp *resource.UpdateResponse) {
|
||||
resp.Diagnostics.AddError("update not supported", "sless_job does not support in-place updates")
|
||||
resp.Diagnostics.AddError("update not supported", "sless_job does not support in-place updates; increment run_id to re-run")
|
||||
}
|
||||
|
||||
func (r *JobResource) Delete(ctx context.Context, req resource.DeleteRequest, resp *resource.DeleteResponse) {
|
||||
@@ -262,19 +325,49 @@ func (r *JobResource) Delete(ctx context.Context, req resource.DeleteRequest, re
|
||||
}
|
||||
}
|
||||
|
||||
// jobToModel конвертирует API-ответ → state модель.
|
||||
// ModifyPlan вычисляет hash директории source_dir на фазе terraform plan.
|
||||
// Если hash изменился относительно state — code_hash в plan будет другим,
|
||||
// что вместе с RequiresReplace на code_hash триггерирует пересоздание джоба.
|
||||
func (r *JobResource) ModifyPlan(ctx context.Context, req resource.ModifyPlanRequest, resp *resource.ModifyPlanResponse) {
|
||||
if req.Plan.Raw.IsNull() {
|
||||
return
|
||||
}
|
||||
var plan JobModel
|
||||
resp.Diagnostics.Append(req.Plan.Get(ctx, &plan)...)
|
||||
if resp.Diagnostics.HasError() {
|
||||
return
|
||||
}
|
||||
if plan.SourceDir.IsNull() || plan.SourceDir.ValueString() == "" {
|
||||
return
|
||||
}
|
||||
_, hash, err := zipDir(plan.SourceDir.ValueString())
|
||||
if err != nil {
|
||||
return // директория может не существовать при первом init — не фейлим plan
|
||||
}
|
||||
plan.CodeHash = types.StringValue(hash)
|
||||
resp.Diagnostics.Append(resp.Plan.Set(ctx, plan)...)
|
||||
}
|
||||
|
||||
// jobToModel конвертирует API-ответ + plan (для локальных полей) → state модель.
|
||||
func jobToModel(plan JobModel, j *client.JobResponse) JobModel {
|
||||
waitTimeoutSec := plan.WaitTimeoutSec
|
||||
if waitTimeoutSec.IsNull() || waitTimeoutSec.IsUnknown() || waitTimeoutSec.ValueInt64() <= 0 {
|
||||
waitTimeoutSec = types.Int64Value(defaultWaitTimeoutSec)
|
||||
waitSec := plan.WaitTimeoutSec
|
||||
if waitSec.IsNull() || waitSec.IsUnknown() || waitSec.ValueInt64() <= 0 {
|
||||
waitSec = types.Int64Value(defaultJobWaitTimeoutSec)
|
||||
}
|
||||
return JobModel{
|
||||
Name: types.StringValue(j.Name),
|
||||
Function: types.StringValue(j.FunctionRef),
|
||||
Runtime: types.StringValue(j.Runtime),
|
||||
Entrypoint: types.StringValue(j.Entrypoint),
|
||||
MemoryMB: plan.MemoryMB,
|
||||
TimeoutSec: plan.TimeoutSec,
|
||||
EnvVars: plan.EnvVars,
|
||||
SourceDir: plan.SourceDir,
|
||||
CodeHash: plan.CodeHash,
|
||||
EventJSON: plan.EventJSON,
|
||||
RunID: types.Int64Value(j.RunID),
|
||||
WaitTimeoutSec: waitTimeoutSec,
|
||||
WaitTimeoutSec: waitSec,
|
||||
Phase: types.StringValue(j.Phase),
|
||||
ImageRef: types.StringValue(j.ImageRef),
|
||||
StartTime: types.StringValue(j.StartTime),
|
||||
CompletionTime: types.StringValue(j.CompletionTime),
|
||||
Message: types.StringValue(j.Message),
|
||||
|
||||
Reference in New Issue
Block a user