Files
sless/controllers/functionjob_controller.go
Naeel a04dfb2d0c fix+docs: FunctionJob label bugfix, job ErrAlreadyExists, python str→text/plain, operator.yaml v0.1.33, progress.md
- controllers/functionjob_controller.go:
  - PodTemplate labels: functionjob=, function= (k8s 1.27+ удалил job-name=)
  - getJobPodOutput принимает labelSelector вместо jobName
  - захват stderr при Failed job; truncateForStatus() helper
- terraform/provider/internal/client/client.go: ErrJobAlreadyExists (409 Conflict)
- terraform/provider/internal/resources/job_resource.go: при конфликте создания — читаем существующий job
- runtimes/python3.11/server.py: str return → text/plain
- internal/builder/context.go: python runtime base image → v0.1.3
- deployments/k8s/operator.yaml: image → v0.1.33
- doc/progress.md: добавлены секции FunctionJob bugfix, str→text/plain, web-console v0.2.0
2026-03-18 17:41:44 +03:00

348 lines
15 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-17 20:00 (bugfix: job-name label удалён в k8s 1.27+, split-brain cached client)
// 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{
ObjectMeta: metav1.ObjectMeta{
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),
Env: append(
append(fnEnvVars(fn), corev1.EnvVar{
Name: "SLESS_EVENT",
Value: eventJSON,
}),
goJobModeEnv(fn.Spec.Runtime)...,
),
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, "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))
`}
}
}
// 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 }
// 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)
}