Files
sless/controllers/functionjob_controller.go

468 lines
20 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Изменено: 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 проверяет наличие образа в registry и либо пропускает сборку,
// либо запускает kaniko Job. Идемпотентность по hash s3Key аналогична service/function.
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
}
imageRef := r.Builder.ImageRef(r.OperatorNamespace, fj.Name, fj.Spec.S3Key)
// Проверяем: образ с этим тегом уже существует в registry?
// Если registry недоступен — requeue, не запускаем сборку.
exists, err := r.Builder.ImageExists(ctx, imageRef)
if err != nil {
return ctrl.Result{RequeueAfter: 10 * time.Second}, fmt.Errorf("check image exists: %w", err)
}
if exists {
logger.Info("image already exists in registry, skipping build", "imageRef", imageRef)
if fj.Annotations == nil {
fj.Annotations = map[string]string{}
}
if err := r.Update(ctx, fj); err != nil {
return ctrl.Result{}, fmt.Errorf("update annotations (cache hit): %w", err)
}
fj.Status.Phase = slessv1alpha1.FunctionJobPhaseBuilding // перейдёт в run сразу
fj.Status.ImageRef = imageRef
fj.Status.Message = "Image restored from registry cache"
if err := r.Status().Update(ctx, fj); err != nil {
return ctrl.Result{}, fmt.Errorf("update status (cache hit): %w", err)
}
return ctrl.Result{Requeue: true}, 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=<name>" (наш лейбл,
// выставляется на 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)
}