- api/v1alpha1/job_types.go: new CRD FunctionJob (Pending/Running/Succeeded/Failed) - controllers/functionjob_controller.go: reconciler creates k8s Job from FunctionRef + EventJSON - zz_generated.deepcopy.go: DeepCopy methods for FunctionJob types - config/crd/bases: generated CRD YAML, applied to cluster - main.go: register FunctionJobReconciler - rbac.yaml: add functionjobs permissions - operator.yaml: v0.1.2 -> v0.1.3 - operator:v0.1.3 deployed and running in cluster
237 lines
9.2 KiB
Go
237 lines
9.2 KiB
Go
// Изменено: 2026-03-07
|
||
// FunctionJobReconciler — контроллер одноразовых запусков функций.
|
||
// При создании FunctionJob:
|
||
// 1. Ждёт пока Function станет Ready
|
||
// 2. Создаёт k8s Job который запускает образ функции с CMD runner
|
||
// 3. Следит за завершением Job → обновляет статус (Succeeded/Failed)
|
||
//
|
||
// Почему отдельный ресурс (не Trigger type=job):
|
||
// Trigger описывает постоянный способ вызова (http endpoint, cron schedule).
|
||
// FunctionJob — разовое событие с отдельным lifecycle и статусом результата.
|
||
|
||
package controllers
|
||
|
||
import (
|
||
"context"
|
||
"fmt"
|
||
|
||
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"
|
||
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"
|
||
)
|
||
|
||
// FunctionJobReconciler reconciles a FunctionJob object
|
||
type FunctionJobReconciler struct {
|
||
client.Client
|
||
Scheme *runtime.Scheme
|
||
RegistrySecret string // имя k8s Secret с docker credentials (для imagePullSecrets)
|
||
}
|
||
|
||
//+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) {
|
||
logger := log.FromContext(ctx)
|
||
|
||
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
|
||
}
|
||
|
||
// Проверяем что 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)
|
||
// Повторный reconcile придёт когда Function изменится
|
||
return ctrl.Result{}, nil
|
||
}
|
||
|
||
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)
|
||
} else if !errors.IsNotFound(err) {
|
||
return ctrl.Result{}, fmt.Errorf("get job: %w", err)
|
||
}
|
||
|
||
// Создаём k8s Job
|
||
// Используем образ функции напрямую, переопределяем CMD чтобы запустить runner
|
||
// вместо server.py/server.js — runner выполняет handle(event) один раз и выходит
|
||
eventJSON := fj.Spec.EventJSON
|
||
if eventJSON == "" {
|
||
eventJSON = "{}"
|
||
}
|
||
|
||
// 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,
|
||
"function": fn.Name,
|
||
},
|
||
},
|
||
Spec: batchv1.JobSpec{
|
||
// Не перезапускать при ошибке — это одноразовый запуск
|
||
BackoffLimit: int32Ptr(0),
|
||
// Автоудаление через 10 мин после завершения — чтобы не засорять кластер
|
||
TTLSecondsAfterFinished: &ttl,
|
||
Template: corev1.PodTemplateSpec{
|
||
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),
|
||
Env: append(
|
||
fnEnvVars(fn),
|
||
corev1.EnvVar{
|
||
Name: "SLESS_EVENT",
|
||
Value: eventJSON,
|
||
},
|
||
),
|
||
Resources: corev1.ResourceRequirements{
|
||
Limits: corev1.ResourceList{
|
||
corev1.ResourceMemory: resource.MustParse(fmt.Sprintf("%dMi", fn.Spec.MemoryMB)),
|
||
},
|
||
},
|
||
},
|
||
},
|
||
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 job for functionjob", "job", jobName, "functionjob", fj.Name)
|
||
return ctrl.Result{}, nil
|
||
}
|
||
|
||
// syncJobStatus читает статус k8s Job и обновляет FunctionJob.Status.
|
||
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
|
||
fj.Status.Message = "completed successfully"
|
||
} else if job.Status.Failed > 0 {
|
||
now := metav1.Now()
|
||
fj.Status.Phase = slessv1alpha1.FunctionJobPhaseFailed
|
||
fj.Status.CompletionTime = &now
|
||
fj.Status.Message = "job failed, check pod logs: kubectl logs -n sless-fn-" + fj.Namespace + " -l functionjob=" + fj.Name
|
||
}
|
||
// Running — ничего не меняем, перечитаем при следующем reconcile
|
||
if err := r.Status().Update(ctx, fj); err != nil {
|
||
return ctrl.Result{}, fmt.Errorf("sync job status: %w", err)
|
||
}
|
||
return ctrl.Result{}, nil
|
||
}
|
||
|
||
// runtimeRunnerCommand возвращает CMD для запуска одноразового runner вместо HTTP-сервера.
|
||
// runner читает env SLESS_EVENT, вызывает handle(event) один раз и завершается.
|
||
func runtimeRunnerCommand(runtime string) []string {
|
||
switch runtime {
|
||
case "nodejs20":
|
||
// inline runner — не требует отдельного файла в образе
|
||
return []string{"node", "-e", `
|
||
const h = require('/app/function/handler.js');
|
||
const event = JSON.parse(process.env.SLESS_EVENT || '{}');
|
||
Promise.resolve(h.handle(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
|
||
spec = importlib.util.spec_from_file_location("handler", "/app/function/handler.py")
|
||
mod = importlib.util.module_from_spec(spec)
|
||
spec.loader.exec_module(mod)
|
||
event = json.loads(os.environ.get("SLESS_EVENT", "{}"))
|
||
result = mod.handle(event)
|
||
print(json.dumps(result))
|
||
`}
|
||
}
|
||
}
|
||
|
||
// fnEnvVars преобразует env vars из FunctionSpec в k8s EnvVar slice.
|
||
func fnEnvVars(fn *slessv1alpha1.Function) []corev1.EnvVar {
|
||
var result []corev1.EnvVar
|
||
for k, v := range fn.Spec.Env {
|
||
result = append(result, corev1.EnvVar{Name: k, Value: v})
|
||
}
|
||
return result
|
||
}
|
||
|
||
func int32Ptr(i int32) *int32 { return &i }
|
||
|
||
// SetupWithManager регистрирует контроллер и настраивает watch на k8s Job.
|
||
func (r *FunctionJobReconciler) SetupWithManager(mgr ctrl.Manager) error {
|
||
return ctrl.NewControllerManagedBy(mgr).
|
||
For(&slessv1alpha1.FunctionJob{}).
|
||
Owns(&batchv1.Job{}).
|
||
Complete(r)
|
||
}
|