feat: operator v0.1.16 — job stdout -> status.Message (feature B)

- FunctionJobReconciler: added KubeClient field (kubernetes.Interface)
- getJobPodOutput(): reads pod logs via typed client after job succeeds
- main.go: inject kubernetes.NewForConfigOrDie into FunctionJobReconciler
- rbac.yaml: add pods/pods/log get/list/watch permissions
- examples/simple-python/: job->function chain demo (Python)
- examples/simple-node/: job->function chain demo (Node.js)

sless_job.X.message now contains the return value of the function
This commit is contained in:
“Naeel”
2026-03-09 14:50:06 +04:00
parent 53d9fa6e74
commit 0aaeb47b3b
19 changed files with 330 additions and 4 deletions
+34 -3
View File
@@ -1,4 +1,4 @@
// Изменено: 2026-03-08 (fix: RequeueAfter для poll статуса Job)
// Изменено: 2026-03-09 (feature B: захват stdout пода Job в status.Message)
// FunctionJobReconciler — контроллер одноразовых запусков функций.
// При создании FunctionJob:
// 1. Ждёт пока Function станет Ready
@@ -12,8 +12,11 @@
package controllers
import (
"bytes"
"context"
"fmt"
"io"
"strings"
"time"
batchv1 "k8s.io/api/batch/v1"
@@ -22,6 +25,7 @@ import (
"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"
@@ -33,7 +37,8 @@ import (
type FunctionJobReconciler struct {
client.Client
Scheme *runtime.Scheme
RegistrySecret string // имя k8s Secret с docker credentials (для imagePullSecrets)
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
@@ -188,7 +193,8 @@ func (r *FunctionJobReconciler) syncJobStatus(ctx context.Context, fj *slessv1al
now := metav1.Now()
fj.Status.Phase = slessv1alpha1.FunctionJobPhaseSucceeded
fj.Status.CompletionTime = &now
fj.Status.Message = "completed successfully"
// Захватываем 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
@@ -259,6 +265,31 @@ func fnEnvVars(fn *slessv1alpha1.Function) []corev1.EnvVar {
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 кросс-неймспейсно не работают.