Files
sless/controllers/functionjob_controller.go
T
“Naeel” 408f58a9e2 fix: stage0 quick fixes (operator v0.1.19)
- TriggerReconciler: RequeueAfter 15s когда Function не Ready
  (ранее зависал без повторного reconcile)
- FunctionJobReconciler: RequeueAfter 15s когда Function не Ready
- UpdateFunction: добавлена валидация runtime/entrypoint/memory_mb
  (ранее мог затереть spec нулями при частичном обновлении)
- CronJob: curlimages/curl:latest → curlimages/curl:8.5.0 (pin version)
- Config: удалён FunctionNamespacePrefix (мёртвое поле, нигде не использовалось)
- Invocations endpoint: возвращает 501 вместо пустого списка
  (SaveInvocation нигде не вызывается — честный ответ клиенту)
- Собран образ naeel/sless-operator:v0.1.19
2026-03-10 17:36:48 +04:00

302 lines
12 KiB
Go
Raw 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-09 (feature B: захват stdout пода Job в status.Message)
// 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 (
"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"
)
// 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)
}
//+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
}
// 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
}
// Проверяем что 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
}
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.
// Если 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, job.Name)
} 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
} 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
}
// runtimeRunnerCommand возвращает CMD для запуска одноразового runner вместо HTTP-сервера.
// runner читает env SLESS_EVENT и SLESS_ENTRYPOINT, вызывает handle(event) один раз и завершается.
func runtimeRunnerCommand(runtime string) []string {
switch runtime {
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))
`}
}
}
// fnEnvVars преобразует env vars из FunctionSpec в k8s EnvVar slice.
// Включает SLESS_ENTRYPOINT чтобы runner.py/runner.js знал какую функцию вызывать.
func fnEnvVars(fn *slessv1alpha1.Function) []corev1.EnvVar {
result := []corev1.EnvVar{
{Name: "SLESS_ENTRYPOINT", Value: fn.Spec.Entrypoint},
}
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 }
// getJobPodOutput находит под созданный Job-ом и возвращает его stdout (trimmed).
// runner.py/runner.js печатают json.dumps(result) в stdout — это и есть return value функции.
// Если под не найден или логи недоступны — возвращает "completed successfully" как fallback.
func getJobPodOutput(ctx context.Context, kube kubernetes.Interface, namespace, jobName string) string {
pods, err := kube.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{
LabelSelector: "job-name=" + jobName,
})
if err != nil || len(pods.Items) == 0 {
return "completed successfully"
}
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)
}