// Изменено: 2026-03-20 (merge sless_function+sless_job: FunctionJobReconciler самодостаточен) // FunctionJobReconciler — контроллер одноразовых запусков функций. // При создании FunctionJob с RunID>0: // 1. Запускает kaniko сборку образа (фаза Building) — больше не зависит от Function CRD // 2. После сборки создаёт k8s Job который запускает образ функции с CMD runner // 3. Следит за завершением Job → обновляет статус (Succeeded/Failed) // // Почему отдельный ресурс (не Trigger type=job): // Trigger описывает постоянный способ вызова (http endpoint, cron schedule). // FunctionJob — разовое событие с отдельным lifecycle и статусом результата. package controllers import ( "bytes" "context" "fmt" "io" "strings" "time" batchv1 "k8s.io/api/batch/v1" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/client-go/kubernetes" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" "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) 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 //+kubebuilder:rbac:groups=sless.kube5s.ru,resources=functionjobs/status,verbs=get;update;patch //+kubebuilder:rbac:groups=sless.kube5s.ru,resources=functionjobs/finalizers,verbs=update //+kubebuilder:rbac:groups=batch,resources=jobs,verbs=get;list;watch;create;update;patch;delete // Reconcile — основной цикл контроллера. func (r *FunctionJobReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) { fj := &slessv1alpha1.FunctionJob{} if err := r.Get(ctx, req.NamespacedName, fj); err != nil { if errors.IsNotFound(err) { return ctrl.Result{}, nil } return ctrl.Result{}, fmt.Errorf("get functionjob: %w", err) } // Уже завершён — ничего не делаем if fj.Status.Phase == slessv1alpha1.FunctionJobPhaseSucceeded || fj.Status.Phase == slessv1alpha1.FunctionJobPhaseFailed { return ctrl.Result{}, nil } // RunID=0 означает "не запускать". // Позволяет создать FunctionJob в кластере, но запустить его явно позже // (увеличив RunID > 0 через Terraform или kubectl patch). if fj.Spec.RunID == 0 { if fj.Status.Phase != "Skipped" { fj.Status.Phase = "Skipped" fj.Status.Message = "run_id=0: set run_id>0 to execute" _ = r.Status().Update(ctx, fj) } return ctrl.Result{}, 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")) existingJob := &batchv1.Job{} if err := r.Get(ctx, client.ObjectKey{Name: jobName, Namespace: deployNS}, existingJob); err == nil { return r.syncJobStatus(ctx, fj, existingJob) } else if !errors.IsNotFound(err) { return ctrl.Result{}, fmt.Errorf("get job: %w", err) } 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) // S3Key пустой — код ещё не загружен (provider делает upload после создания CRD). // Ждём: через 5 секунд provider успеет выполнить UploadJobCode → S3Key заполнится. if fj.Spec.S3Key == "" { logger.Info("s3key empty, waiting for code upload", "job", fj.Name) return ctrl.Result{RequeueAfter: 5 * time.Second}, nil } 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 мин после завершения job := &batchv1.Job{ ObjectMeta: metav1.ObjectMeta{ Name: jobName, Namespace: deployNS, Labels: map[string]string{ "managed-by": "sless", "functionjob": fj.Name, }, }, Spec: batchv1.JobSpec{ // Не перезапускать при ошибке — это одноразовый запуск BackoffLimit: int32Ptr(0), // Автоудаление через 10 мин после завершения — чтобы не засорять кластер TTLSecondsAfterFinished: &ttl, Template: corev1.PodTemplateSpec{ ObjectMeta: metav1.ObjectMeta{ Labels: map[string]string{ "managed-by": "sless", "functionjob": fj.Name, }, }, Spec: corev1.PodSpec{ RestartPolicy: corev1.RestartPolicyNever, Containers: []corev1.Container{ { Name: "runner", Image: fj.Status.ImageRef, Command: runtimeRunnerCommand(fj.Spec.Runtime), Env: append( append(fjEnvVars(fj), corev1.EnvVar{ Name: "SLESS_EVENT", Value: eventJSON, }), goJobModeEnv(fj.Spec.Runtime)..., ), Resources: corev1.ResourceRequirements{ Limits: corev1.ResourceList{ corev1.ResourceMemory: resource.MustParse(fmt.Sprintf("%dMi", memMB)), }, }, }, }, ImagePullSecrets: []corev1.LocalObjectReference{ {Name: r.RegistrySecret}, }, }, }, }, } if err := r.Create(ctx, job); err != nil { return ctrl.Result{}, fmt.Errorf("create job: %w", err) } now := metav1.Now() fj.Status.Phase = slessv1alpha1.FunctionJobPhaseRunning fj.Status.JobName = jobName fj.Status.StartTime = &now fj.Status.Message = "" if err := r.Status().Update(ctx, fj); err != nil { return ctrl.Result{}, fmt.Errorf("update functionjob status: %w", err) } 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). func (r *FunctionJobReconciler) syncJobStatus(ctx context.Context, fj *slessv1alpha1.FunctionJob, job *batchv1.Job) (ctrl.Result, error) { if job.Status.Succeeded > 0 { now := metav1.Now() fj.Status.Phase = slessv1alpha1.FunctionJobPhaseSucceeded fj.Status.CompletionTime = &now // Захватываем stdout пода — это return value функции (runner делает print(json.dumps(result))) fj.Status.Message = getJobPodOutput(ctx, r.KubeClient, job.Namespace, "functionjob="+fj.Name) } else if job.Status.Failed > 0 { now := metav1.Now() fj.Status.Phase = slessv1alpha1.FunctionJobPhaseFailed fj.Status.CompletionTime = &now // Захватываем логи по нашему лейблу functionjob= (работает во всех версиях k8s). // Устаревший job-name= удалён в k8s 1.27+, batch.kubernetes.io/job-name= — только с 1.27. podOutput := strings.TrimSpace(getJobPodOutput(ctx, r.KubeClient, job.Namespace, "functionjob="+fj.Name)) if podOutput == "" || podOutput == "completed successfully" { fj.Status.Message = "job failed, check pod logs: kubectl logs -n " + job.Namespace + " -l functionjob=" + fj.Name } else { fj.Status.Message = "job failed: " + truncateForStatus(podOutput, 2000) } } else { // Job ещё выполняется — перечитаем через 5 секунд if err := r.Status().Update(ctx, fj); err != nil { return ctrl.Result{}, fmt.Errorf("sync job status: %w", err) } return ctrl.Result{RequeueAfter: 5 * time.Second}, nil } if err := r.Status().Update(ctx, fj); err != nil { return ctrl.Result{}, fmt.Errorf("sync job status: %w", err) } return ctrl.Result{}, nil } // truncateForStatus ограничивает длину текста для безопасной записи в status.message. func truncateForStatus(message string, maxLen int) string { if len(message) <= maxLen { return message } if maxLen <= 3 { return message[:maxLen] } return message[:maxLen-3] + "..." } // runtimeRunnerCommand возвращает CMD для запуска одноразового runner вместо HTTP-сервера. // runner читает env SLESS_EVENT и SLESS_ENTRYPOINT, вызывает handle(event) один раз и завершается. func runtimeRunnerCommand(runtime string) []string { switch runtime { case "go1.23": // Go runtime: SLESS_MODE=job заставляет server читать SLESS_EVENT и выйти. // CMD остаётся как в образе (/server), переопределяем через Env. // Передаём пустую команду — используется CMD из образа (/server). // SLESS_MODE=job добавляется через Env в fnEnvVars. return nil // nil = использовать CMD из образа; SLESS_MODE=job в Env case "nodejs20": // inline runner — не требует отдельного файла в образе. // SLESS_ENTRYPOINT="module.func": module=имя файла, func=экспортируемая функция return []string{"node", "-e", ` const ep = process.env.SLESS_ENTRYPOINT || 'handler.handle'; const dot = ep.lastIndexOf('.'); const mod = ep.slice(0, dot >= 0 ? dot : ep.length); const fn = dot >= 0 ? ep.slice(dot + 1) : 'handle'; const h = require('/app/function/' + mod); const event = JSON.parse(process.env.SLESS_EVENT || '{}'); Promise.resolve(h[fn](event)).then(r => { console.log(JSON.stringify(r)); process.exit(0); }).catch(e => { console.error(e.message); process.exit(1); });`} default: // python3.11 return []string{"python3", "-c", ` import os, json, importlib.util ep = os.environ.get("SLESS_ENTRYPOINT", "handler.handle") dot = ep.rfind(".") mod_name = ep[:dot] if dot >= 0 else ep fn_name = ep[dot+1:] if dot >= 0 else "handle" spec = importlib.util.spec_from_file_location(mod_name, "/app/function/" + mod_name + ".py") mod = importlib.util.module_from_spec(spec) spec.loader.exec_module(mod) event = json.loads(os.environ.get("SLESS_EVENT", "{}")) result = getattr(mod, fn_name)(event) print(json.dumps(result)) `} } } // fjEnvVars формирует k8s EnvVar из полей FunctionJobSpec. // SLESS_ENTRYPOINT сообщает runner'у какую функцию вызывать. func fjEnvVars(fj *slessv1alpha1.FunctionJob) []corev1.EnvVar { result := []corev1.EnvVar{ {Name: "SLESS_ENTRYPOINT", Value: fj.Spec.Entrypoint}, } for k, v := range fj.Spec.Env { result = append(result, corev1.EnvVar{Name: k, Value: v}) } return result } func int32Ptr(i int32) *int32 { return &i } // goJobModeEnv возвращает SLESS_MODE=job для Go runtime — сигнал /server выполниться разово и выйти. // Для Python/Node runner задаётся через Command, для Go — через env (CMD /server общий). func goJobModeEnv(runtime string) []corev1.EnvVar { if runtime == "go1.23" { return []corev1.EnvVar{{Name: "SLESS_MODE", Value: "job"}} } return nil } // getJobPodOutput находит под по labelSelector и возвращает его stdout+stderr (trimmed). // runner.py/runner.js печатают json.dumps(result) в stdout — return value функции. // Исключения/трейсбэки Python/Node пишут в stderr — поэтому собираем оба потока. // Если под не найден или логи недоступны — возвращает "completed successfully" как fallback. // labelSelector передаётся снаружи — вызывающий код использует "functionjob=" (наш лейбл, // выставляется на PodTemplate контроллером и не зависит от версии k8s). // НЕ использовать "job-name=" — этот встроенный лейбл удалён в k8s 1.27+ (у нас 1.34.1). func getJobPodOutput(ctx context.Context, kube kubernetes.Interface, namespace, labelSelector string) string { pods, err := kube.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{ LabelSelector: labelSelector, }) if err != nil || len(pods.Items) == 0 { return "completed successfully" } // Stdout: true, Stderr: true — собираем оба потока. // Python исключения идут в stderr, runner.py пишет результат в stdout. req := kube.CoreV1().Pods(namespace).GetLogs(pods.Items[0].Name, &corev1.PodLogOptions{}) stream, err := req.Stream(ctx) if err != nil { return "completed successfully" } defer stream.Close() buf := new(bytes.Buffer) _, _ = io.Copy(buf, stream) out := strings.TrimSpace(buf.String()) if out == "" { return "completed successfully" } return out } // SetupWithManager регистрирует контроллер. // Owns(&batchv1.Job{}) намеренно убрано: Job создаётся в другом namespace // (sless-fn-*), где OwnerReference кросс-неймспейсно не работают. // Вместо этого используется RequeueAfter-polling в syncJobStatus. func (r *FunctionJobReconciler) SetupWithManager(mgr ctrl.Manager) error { return ctrl.NewControllerManagedBy(mgr). For(&slessv1alpha1.FunctionJob{}). Complete(r) }