diff --git a/.gitignore b/.gitignore index 4741072..9ae5e6b 100644 --- a/.gitignore +++ b/.gitignore @@ -65,3 +65,7 @@ plan.out sless-plan examples/.git event-dispatcher + +# build artifacts +/sless +examples/POSTGRES/stress_log*.txt diff --git a/api/v1alpha1/job_types.go b/api/v1alpha1/job_types.go index eec8504..4e4ac13 100644 --- a/api/v1alpha1/job_types.go +++ b/api/v1alpha1/job_types.go @@ -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` diff --git a/api/v1alpha1/zz_generated.deepcopy.go b/api/v1alpha1/zz_generated.deepcopy.go index 8153617..f26f0a8 100644 --- a/api/v1alpha1/zz_generated.deepcopy.go +++ b/api/v1alpha1/zz_generated.deepcopy.go @@ -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. diff --git a/controllers/functionjob_controller.go b/controllers/functionjob_controller.go index af94795..d9b77dc 100644 --- a/controllers/functionjob_controller.go +++ b/controllers/functionjob_controller.go @@ -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 diff --git a/doc/decisions/log.md b/doc/decisions/log.md index 5eab90d..5534970 100644 --- a/doc/decisions/log.md +++ b/doc/decisions/log.md @@ -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 ### Контекст diff --git a/doc/progress.md b/doc/progress.md index eba2a35..a911fde 100644 --- a/doc/progress.md +++ b/doc/progress.md @@ -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 ### Что произошло diff --git a/examples/POSTGRES/functions.tf b/examples/POSTGRES/functions.tf index cf50036..e2f553c 100644 --- a/examples/POSTGRES/functions.tf +++ b/examples/POSTGRES/functions.tf @@ -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//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//pg-table-reader -# https://sless.kube5s.ru/fn//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 } - -*/ diff --git a/internal/api/handler/jobs.go b/internal/api/handler/jobs.go index 30010c3..251b512 100644 --- a/internal/api/handler/jobs.go +++ b/internal/api/handler/jobs.go @@ -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", + }) +} diff --git a/internal/api/router.go b/internal/api/router.go index 4f43259..17d9bd1 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -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. diff --git a/main.go b/main.go index 9fe1697..d46e65c 100644 --- a/main.go +++ b/main.go @@ -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) diff --git a/terraform/provider/internal/client/client.go b/terraform/provider/internal/client/client.go index 058ee18..920e7c4 100644 --- a/terraform/provider/internal/client/client.go +++ b/terraform/provider/internal/client/client.go @@ -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"` diff --git a/terraform/provider/internal/resources/job_resource.go b/terraform/provider/internal/resources/job_resource.go index ad31fb6..33b188e 100644 --- a/terraform/provider/internal/resources/job_resource.go +++ b/terraform/provider/internal/resources/job_resource.go @@ -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),