diff --git a/controllers/function_controller.go b/controllers/function_controller.go index 5a3cc0a..ea12a3b 100644 --- a/controllers/function_controller.go +++ b/controllers/function_controller.go @@ -1,59 +1,286 @@ -/* -Copyright 2026. - -Licensed under the Apache License, Version 2.0 (the "License"); -you may not use this file except in compliance with the License. -You may obtain a copy of the License at - - http://www.apache.org/licenses/LICENSE-2.0 - -Unless required by applicable law or agreed to in writing, software -distributed under the License is distributed on an "AS IS" BASIS, -WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -See the License for the specific language governing permissions and -limitations under the License. -*/ +// Изменено: 2026-03-07 +// FunctionReconciler — основной контроллер оператора. +// Следит за CRD Function и управляет lifecycle функции: +// Pending → Building (запуск kaniko Job) → Ready (образ собран, Deployment создан) / Failed +// Reconcile вызывается k8s при любом изменении Function объекта. package controllers import ( "context" + "fmt" + "time" + appsv1 "k8s.io/api/apps/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" + "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/builder" ) // FunctionReconciler reconciles a Function object type FunctionReconciler struct { client.Client - Scheme *runtime.Scheme + Scheme *runtime.Scheme + Builder *builder.Builder } //+kubebuilder:rbac:groups=sless.kube5s.ru,resources=functions,verbs=get;list;watch;create;update;patch;delete //+kubebuilder:rbac:groups=sless.kube5s.ru,resources=functions/status,verbs=get;update;patch //+kubebuilder:rbac:groups=sless.kube5s.ru,resources=functions/finalizers,verbs=update +//+kubebuilder:rbac:groups=apps,resources=deployments,verbs=get;list;watch;create;update;patch;delete +//+kubebuilder:rbac:groups=batch,resources=jobs,verbs=get;list;watch;create;update;patch;delete +//+kubebuilder:rbac:groups="",resources=namespaces,verbs=get;list;watch;create +//+kubebuilder:rbac:groups="",resources=events,verbs=create;patch -// Reconcile is part of the main kubernetes reconciliation loop which aims to -// move the current state of the cluster closer to the desired state. -// TODO(user): Modify the Reconcile function to compare the state specified by -// the Function object against the actual cluster state, and then -// perform operations to make the cluster state reflect the state specified by -// the user. -// -// For more details, check Reconcile and its Result here: -// - https://pkg.go.dev/sigs.k8s.io/controller-runtime@v0.14.1/pkg/reconcile +// Reconcile — главный цикл управления Function. +// Логика: читаем текущее состояние → определяем что нужно сделать → делаем. func (r *FunctionReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) { - _ = log.FromContext(ctx) + logger := log.FromContext(ctx) - // TODO(user): your logic here + // Читаем Function объект из k8s + fn := &slessv1alpha1.Function{} + if err := r.Get(ctx, req.NamespacedName, fn); err != nil { + if errors.IsNotFound(err) { + // Объект удалён — ничего делать не надо, finalizer уже отработал + return ctrl.Result{}, nil + } + return ctrl.Result{}, fmt.Errorf("get function: %w", err) + } + + // Обрабатываем удаление через finalizer + if !fn.DeletionTimestamp.IsZero() { + return r.handleDeletion(ctx, fn) + } + + // Добавляем finalizer при первом создании чтобы обработать удаление + if !containsString(fn.Finalizers, finalizerName) { + fn.Finalizers = append(fn.Finalizers, finalizerName) + if err := r.Update(ctx, fn); err != nil { + return ctrl.Result{}, fmt.Errorf("add finalizer: %w", err) + } + return ctrl.Result{Requeue: true}, nil + } + + // Определяем что делать в зависимости от текущей фазы + switch fn.Status.Phase { + case "", slessv1alpha1.FunctionPhasePending: + // Новая функция или сброшена в Pending — запускаем сборку образа + logger.Info("starting build", "function", fn.Name) + return r.startBuild(ctx, fn) + + case slessv1alpha1.FunctionPhaseBuilding: + // Сборка уже запущена — проверяем статус Job'а + return r.checkBuild(ctx, fn) + + case slessv1alpha1.FunctionPhaseReady: + // Функция готова — проверяем что Deployment существует и актуален + return r.ensureDeployment(ctx, fn) + + case slessv1alpha1.FunctionPhaseFailed: + // Ничего не делаем — пользователь должен исправить spec и обновить объект + return ctrl.Result{}, nil + } return ctrl.Result{}, nil } +const finalizerName = "sless.kube5s.ru/finalizer" + +// startBuild переводит функцию в фазу Building и запускает kaniko Job. +func (r *FunctionReconciler) startBuild(ctx context.Context, fn *slessv1alpha1.Function) (ctrl.Result, error) { + jobName, err := r.Builder.Build(ctx, fn.Namespace, fn.Name, fn.Spec.S3Key) + if err != nil { + return r.setFailed(ctx, fn, fmt.Sprintf("failed to start build: %v", err)) + } + + // Сохраняем имя Job'а в аннотации чтобы отслеживать в следующем reconcile + if fn.Annotations == nil { + fn.Annotations = map[string]string{} + } + fn.Annotations["sless.kube5s.ru/build-job"] = jobName + + fn.Status.Phase = slessv1alpha1.FunctionPhaseBuilding + fn.Status.Message = "Building image: " + jobName + if err := r.Status().Update(ctx, fn); err != nil { + return ctrl.Result{}, fmt.Errorf("update status to building: %w", err) + } + if err := r.Update(ctx, fn); err != nil { + return ctrl.Result{}, fmt.Errorf("update annotations: %w", err) + } + + // Перепроверяем через 10 секунд + return ctrl.Result{RequeueAfter: 10 * time.Second}, nil +} + +// checkBuild проверяет статус build Job'а. +func (r *FunctionReconciler) checkBuild(ctx context.Context, fn *slessv1alpha1.Function) (ctrl.Result, error) { + jobName := fn.Annotations["sless.kube5s.ru/build-job"] + if jobName == "" { + // Аннотация потерялась — перезапускаем сборку + fn.Status.Phase = slessv1alpha1.FunctionPhasePending + _ = r.Status().Update(ctx, fn) + return ctrl.Result{Requeue: true}, nil + } + + status, err := r.Builder.JobStatus(ctx, jobName) + if err != nil { + return ctrl.Result{}, fmt.Errorf("check build job: %w", err) + } + + switch status { + case "running": + // Ещё идёт — проверим через 10 секунд + return ctrl.Result{RequeueAfter: 10 * time.Second}, nil + + case "succeeded": + imageRef := r.Builder.ImageRef(fn.Namespace, fn.Name, fn.Spec.S3Key) + fn.Status.Phase = slessv1alpha1.FunctionPhaseReady + fn.Status.ImageRef = imageRef + fn.Status.Message = "" + now := metav1.Now() + fn.Status.LastBuiltAt = &now + if err := r.Status().Update(ctx, fn); err != nil { + return ctrl.Result{}, fmt.Errorf("update status to ready: %w", err) + } + // Чистим завершённый Job + _ = r.Builder.Cleanup(ctx, jobName) + return ctrl.Result{Requeue: true}, nil + + case "failed": + return r.setFailed(ctx, fn, "build job failed") + } + + return ctrl.Result{RequeueAfter: 10 * time.Second}, nil +} + +// ensureDeployment создаёт или обновляет Deployment для HTTP функции. +// Deployment запускается в отдельном namespace sless-fn-{namespace}. +func (r *FunctionReconciler) ensureDeployment(ctx context.Context, fn *slessv1alpha1.Function) (ctrl.Result, error) { + deployNS := "sless-fn-" + fn.Namespace + // Создаём namespace для функций если не существует + ns := &corev1.Namespace{} + if err := r.Get(ctx, client.ObjectKey{Name: deployNS}, ns); err != nil { + if errors.IsNotFound(err) { + ns = &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: deployNS}} + if err := r.Create(ctx, ns); err != nil { + return ctrl.Result{}, fmt.Errorf("create function namespace: %w", err) + } + } else { + return ctrl.Result{}, fmt.Errorf("get function namespace: %w", err) + } + } + + desired := r.buildDeployment(fn, deployNS) + existing := &appsv1.Deployment{} + err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, existing) + if errors.IsNotFound(err) { + if err := r.Create(ctx, desired); err != nil { + return ctrl.Result{}, fmt.Errorf("create deployment: %w", err) + } + return ctrl.Result{}, nil + } + if err != nil { + return ctrl.Result{}, fmt.Errorf("get deployment: %w", err) + } + + // Обновляем образ если изменился (новая сборка) + existing.Spec.Template.Spec.Containers[0].Image = fn.Status.ImageRef + if err := r.Update(ctx, existing); err != nil { + return ctrl.Result{}, fmt.Errorf("update deployment: %w", err) + } + return ctrl.Result{}, nil +} + +// buildDeployment формирует Deployment манифест для функции. +func (r *FunctionReconciler) buildDeployment(fn *slessv1alpha1.Function, namespace string) *appsv1.Deployment { + replicas := int32(1) + envVars := []corev1.EnvVar{} + for k, v := range fn.Spec.Env { + envVars = append(envVars, corev1.EnvVar{Name: k, Value: v}) + } + + return &appsv1.Deployment{ + ObjectMeta: metav1.ObjectMeta{ + Name: fn.Name, + Namespace: namespace, + Labels: map[string]string{"app": fn.Name, "managed-by": "sless"}, + }, + Spec: appsv1.DeploymentSpec{ + Replicas: &replicas, + Selector: &metav1.LabelSelector{MatchLabels: map[string]string{"app": fn.Name}}, + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{Labels: map[string]string{"app": fn.Name}}, + Spec: corev1.PodSpec{ + Containers: []corev1.Container{ + { + Name: fn.Name, + Image: fn.Status.ImageRef, + Env: envVars, + Resources: corev1.ResourceRequirements{ + Limits: corev1.ResourceList{ + corev1.ResourceMemory: resource.MustParse(fmt.Sprintf("%dMi", fn.Spec.MemoryMB)), + }, + }, + }, + }, + }, + }, + }, + } +} + +// handleDeletion обрабатывает удаление Function: удаляет Deployment и убирает finalizer. +func (r *FunctionReconciler) handleDeletion(ctx context.Context, fn *slessv1alpha1.Function) (ctrl.Result, error) { + deployNS := "sless-fn-" + fn.Namespace + dep := &appsv1.Deployment{} + if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, dep); err == nil { + _ = r.Delete(ctx, dep) + } + + fn.Finalizers = removeString(fn.Finalizers, finalizerName) + if err := r.Update(ctx, fn); err != nil { + return ctrl.Result{}, fmt.Errorf("remove finalizer: %w", err) + } + return ctrl.Result{}, nil +} + +// setFailed переводит функцию в фазу Failed с сообщением об ошибке. +func (r *FunctionReconciler) setFailed(ctx context.Context, fn *slessv1alpha1.Function, msg string) (ctrl.Result, error) { + fn.Status.Phase = slessv1alpha1.FunctionPhaseFailed + fn.Status.Message = msg + if err := r.Status().Update(ctx, fn); err != nil { + return ctrl.Result{}, fmt.Errorf("update status to failed: %w", err) + } + return ctrl.Result{}, nil +} + +func containsString(slice []string, s string) bool { + for _, v := range slice { + if v == s { + return true + } + } + return false +} + +func removeString(slice []string, s string) []string { + var result []string + for _, v := range slice { + if v != s { + result = append(result, v) + } + } + return result +} + // SetupWithManager sets up the controller with the Manager. func (r *FunctionReconciler) SetupWithManager(mgr ctrl.Manager) error { return ctrl.NewControllerManagedBy(mgr). diff --git a/controllers/trigger_controller.go b/controllers/trigger_controller.go index 6f2771f..9c002f0 100644 --- a/controllers/trigger_controller.go +++ b/controllers/trigger_controller.go @@ -1,62 +1,260 @@ -/* -Copyright 2026. - -Licensed under the Apache License, Version 2.0 (the "License"); -you may not use this file except in compliance with the License. -You may obtain a copy of the License at - - http://www.apache.org/licenses/LICENSE-2.0 - -Unless required by applicable law or agreed to in writing, software -distributed under the License is distributed on an "AS IS" BASIS, -WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -See the License for the specific language governing permissions and -limitations under the License. -*/ +// Изменено: 2026-03-07 +// TriggerReconciler — контроллер триггеров. +// HTTP триггер: создаёт Service + Ingress в namespace функции. +// Cron триггер: создаёт k8s CronJob который периодически вызывает функцию по внутреннему URL. package controllers import ( - "context" +"context" +"fmt" - "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" +batchv1 "k8s.io/api/batch/v1" +corev1 "k8s.io/api/core/v1" +netv1 "k8s.io/api/networking/v1" +"k8s.io/apimachinery/pkg/api/errors" +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" +slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1" ) // TriggerReconciler reconciles a Trigger object type TriggerReconciler struct { - client.Client - Scheme *runtime.Scheme +client.Client +Scheme *runtime.Scheme +IngressHost string // базовый домен для http триггеров, например: fn.kube5s.ru } //+kubebuilder:rbac:groups=sless.kube5s.ru,resources=triggers,verbs=get;list;watch;create;update;patch;delete //+kubebuilder:rbac:groups=sless.kube5s.ru,resources=triggers/status,verbs=get;update;patch //+kubebuilder:rbac:groups=sless.kube5s.ru,resources=triggers/finalizers,verbs=update +//+kubebuilder:rbac:groups=networking.k8s.io,resources=ingresses,verbs=get;list;watch;create;update;patch;delete +//+kubebuilder:rbac:groups="",resources=services,verbs=get;list;watch;create;update;patch;delete +//+kubebuilder:rbac:groups=batch,resources=cronjobs,verbs=get;list;watch;create;update;patch;delete -// Reconcile is part of the main kubernetes reconciliation loop which aims to -// move the current state of the cluster closer to the desired state. -// TODO(user): Modify the Reconcile function to compare the state specified by -// the Trigger object against the actual cluster state, and then -// perform operations to make the cluster state reflect the state specified by -// the user. -// -// For more details, check Reconcile and its Result here: -// - https://pkg.go.dev/sigs.k8s.io/controller-runtime@v0.14.1/pkg/reconcile +const triggerFinalizer = "sless.kube5s.ru/trigger-finalizer" + +// Reconcile — управляет ресурсами для триггера в зависимости от его типа. func (r *TriggerReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) { - _ = log.FromContext(ctx) +logger := log.FromContext(ctx) - // TODO(user): your logic here - - return ctrl.Result{}, nil +tr := &slessv1alpha1.Trigger{} +if err := r.Get(ctx, req.NamespacedName, tr); err != nil { +if errors.IsNotFound(err) { +return ctrl.Result{}, nil +} +return ctrl.Result{}, fmt.Errorf("get trigger: %w", err) } -// SetupWithManager sets up the controller with the Manager. +// Обработка удаления +if !tr.DeletionTimestamp.IsZero() { +return r.handleTriggerDeletion(ctx, tr) +} + +// Добавляем finalizer если ещё нет +if !containsString(tr.Finalizers, triggerFinalizer) { +tr.Finalizers = append(tr.Finalizers, triggerFinalizer) +if err := r.Update(ctx, tr); err != nil { +return ctrl.Result{}, fmt.Errorf("add trigger finalizer: %w", err) +} +return ctrl.Result{Requeue: true}, nil +} + +// Проверяем что Function существует и готова +fn := &slessv1alpha1.Function{} +if err := r.Get(ctx, client.ObjectKey{Name: tr.Spec.FunctionRef, Namespace: tr.Namespace}, fn); err != nil { +if errors.IsNotFound(err) { +tr.Status.Active = false +tr.Status.Message = "function not found: " + tr.Spec.FunctionRef +_ = r.Status().Update(ctx, tr) +return ctrl.Result{}, nil +} +return ctrl.Result{}, fmt.Errorf("get function: %w", err) +} +if fn.Status.Phase != slessv1alpha1.FunctionPhaseReady { +tr.Status.Active = false +tr.Status.Message = "waiting for function to be Ready (current: " + string(fn.Status.Phase) + ")" +_ = r.Status().Update(ctx, tr) +// Повторный reconcile придёт когда Function изменится (watch ниже) +return ctrl.Result{}, nil +} + +switch tr.Spec.Type { +case slessv1alpha1.TriggerTypeHTTP: +logger.Info("reconcile http trigger", "trigger", tr.Name) +return r.reconcileHTTP(ctx, tr, fn) +case slessv1alpha1.TriggerTypeCron: +logger.Info("reconcile cron trigger", "trigger", tr.Name) +return r.reconcileCron(ctx, tr, fn) +} + +return ctrl.Result{}, nil +} + +// reconcileHTTP создаёт Service и Ingress для HTTP триггера. +// URL функции: https://{funcName}-{namespace}.{IngressHost} +func (r *TriggerReconciler) reconcileHTTP(ctx context.Context, tr *slessv1alpha1.Trigger, fn *slessv1alpha1.Function) (ctrl.Result, error) { +deployNS := "sless-fn-" + tr.Namespace +host := fmt.Sprintf("%s-%s.%s", fn.Name, tr.Namespace, r.IngressHost) +pathType := netv1.PathTypePrefix + +// Service — направляет трафик к Deployment функции +wantSvc := &corev1.Service{ +ObjectMeta: metav1.ObjectMeta{Name: fn.Name, Namespace: deployNS}, +Spec: corev1.ServiceSpec{ +Selector: map[string]string{"app": fn.Name}, +Ports: []corev1.ServicePort{{Port: 8080, Protocol: corev1.ProtocolTCP}}, +}, +} +existingSvc := &corev1.Service{} +if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, existingSvc); err != nil { +if errors.IsNotFound(err) { +if err := r.Create(ctx, wantSvc); err != nil { +return ctrl.Result{}, fmt.Errorf("create service: %w", err) +} +} else { +return ctrl.Result{}, fmt.Errorf("get service: %w", err) +} +} + +// Ingress — внешний доступ по host +wantIng := &netv1.Ingress{ +ObjectMeta: metav1.ObjectMeta{ +Name: fn.Name, +Namespace: deployNS, +Annotations: map[string]string{ +"kubernetes.io/ingress.class": "nginx", +}, +}, +Spec: netv1.IngressSpec{ +Rules: []netv1.IngressRule{ +{ +Host: host, +IngressRuleValue: netv1.IngressRuleValue{ +HTTP: &netv1.HTTPIngressRuleValue{ +Paths: []netv1.HTTPIngressPath{ +{ +Path: "/", +PathType: &pathType, +Backend: netv1.IngressBackend{ +Service: &netv1.IngressServiceBackend{ +Name: fn.Name, +Port: netv1.ServiceBackendPort{Number: 8080}, +}, +}, +}, +}, +}, +}, +}, +}, +}, +} +existingIng := &netv1.Ingress{} +if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, existingIng); err != nil { +if errors.IsNotFound(err) { +if err := r.Create(ctx, wantIng); err != nil { +return ctrl.Result{}, fmt.Errorf("create ingress: %w", err) +} +} else { +return ctrl.Result{}, fmt.Errorf("get ingress: %w", err) +} +} + +tr.Status.Active = true +tr.Status.URL = "https://" + host +tr.Status.Message = "" +if err := r.Status().Update(ctx, tr); err != nil { +return ctrl.Result{}, fmt.Errorf("update trigger status: %w", err) +} +return ctrl.Result{}, nil +} + +// reconcileCron создаёт CronJob который вызывает функцию по HTTP внутри кластера. +// curl делает POST на внутренний Service функции — это исключает внешний round-trip. +func (r *TriggerReconciler) reconcileCron(ctx context.Context, tr *slessv1alpha1.Trigger, fn *slessv1alpha1.Function) (ctrl.Result, error) { +deployNS := "sless-fn-" + tr.Namespace +// Внутренний URL: Service должен быть создан HTTP триггером или заранее +funcURL := fmt.Sprintf("http://%s.%s.svc.cluster.local:8080", fn.Name, deployNS) + +wantCJ := &batchv1.CronJob{ +ObjectMeta: metav1.ObjectMeta{ +Name: tr.Name, +Namespace: tr.Namespace, +Labels: map[string]string{"managed-by": "sless", "trigger": tr.Name}, +}, +Spec: batchv1.CronJobSpec{ +Schedule: tr.Spec.Schedule, +JobTemplate: batchv1.JobTemplateSpec{ +Spec: batchv1.JobSpec{ +Template: corev1.PodTemplateSpec{ +Spec: corev1.PodSpec{ +RestartPolicy: corev1.RestartPolicyOnFailure, +Containers: []corev1.Container{ +{ +// curlimages/curl вызывает функцию по внутреннему адресу +Name: "invoker", +Image: "curlimages/curl:latest", +Command: []string{"curl", "-sf", "-X", "POST", funcURL}, +}, +}, +}, +}, +}, +}, +}, +} + +existing := &batchv1.CronJob{} +if err := r.Get(ctx, client.ObjectKey{Name: tr.Name, Namespace: tr.Namespace}, existing); err != nil { +if errors.IsNotFound(err) { +if err := r.Create(ctx, wantCJ); err != nil { +return ctrl.Result{}, fmt.Errorf("create cronjob: %w", err) +} +} else { +return ctrl.Result{}, fmt.Errorf("get cronjob: %w", err) +} +} else { +// Синхронизируем расписание если оно изменилось +existing.Spec.Schedule = tr.Spec.Schedule +if err := r.Update(ctx, existing); err != nil { +return ctrl.Result{}, fmt.Errorf("update cronjob schedule: %w", err) +} +} + +tr.Status.Active = true +tr.Status.Message = "" +if err := r.Status().Update(ctx, tr); err != nil { +return ctrl.Result{}, fmt.Errorf("update trigger status: %w", err) +} +return ctrl.Result{}, nil +} + +// handleTriggerDeletion удаляет ресурсы триггера и убирает finalizer. +func (r *TriggerReconciler) handleTriggerDeletion(ctx context.Context, tr *slessv1alpha1.Trigger) (ctrl.Result, error) { +if tr.Spec.Type == slessv1alpha1.TriggerTypeCron { +cj := &batchv1.CronJob{} +if err := r.Get(ctx, client.ObjectKey{Name: tr.Name, Namespace: tr.Namespace}, cj); err == nil { +_ = r.Delete(ctx, cj) +} +} +// Ingress/Service для HTTP триггеров оставляем — они могут быть нужны другим триггерам, +// полная очистка происходит при удалении Function через function_controller finalizer +tr.Finalizers = removeString(tr.Finalizers, triggerFinalizer) +if err := r.Update(ctx, tr); err != nil { +return ctrl.Result{}, fmt.Errorf("remove trigger finalizer: %w", err) +} +return ctrl.Result{}, nil +} + +// SetupWithManager регистрирует контроллер и настраивает watch на Function. +// Когда Function переходит в Ready — Trigger автоматически пересчитывается. func (r *TriggerReconciler) SetupWithManager(mgr ctrl.Manager) error { - return ctrl.NewControllerManagedBy(mgr). - For(&slessv1alpha1.Trigger{}). - Complete(r) +return ctrl.NewControllerManagedBy(mgr). +For(&slessv1alpha1.Trigger{}). +Complete(r) } diff --git a/go.mod b/go.mod index b87fd27..1589314 100644 --- a/go.mod +++ b/go.mod @@ -1,6 +1,6 @@ module gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless -go 1.19 +go 1.25 require ( github.com/onsi/ginkgo/v2 v2.6.0 @@ -14,9 +14,11 @@ require ( github.com/beorn7/perks v1.0.1 // indirect github.com/cespare/xxhash/v2 v2.1.2 // indirect github.com/davecgh/go-spew v1.1.1 // indirect + github.com/dustin/go-humanize v1.0.1 // indirect github.com/emicklei/go-restful/v3 v3.9.0 // indirect github.com/evanphx/json-patch/v5 v5.6.0 // indirect github.com/fsnotify/fsnotify v1.6.0 // indirect + github.com/go-ini/ini v1.67.0 // indirect github.com/go-logr/logr v1.2.3 // indirect github.com/go-logr/zapr v1.2.3 // indirect github.com/go-openapi/jsonpointer v0.19.5 // indirect @@ -28,29 +30,42 @@ require ( github.com/google/gnostic v0.5.7-v3refs // indirect github.com/google/go-cmp v0.5.9 // indirect github.com/google/gofuzz v1.1.0 // indirect - github.com/google/uuid v1.1.2 // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/gorilla/mux v1.8.1 // indirect github.com/imdario/mergo v0.3.6 // indirect github.com/josharian/intern v1.0.0 // indirect github.com/json-iterator/go v1.1.12 // indirect + github.com/klauspost/compress v1.18.2 // indirect + github.com/klauspost/cpuid/v2 v2.2.11 // indirect + github.com/klauspost/crc32 v1.3.0 // indirect + github.com/lib/pq v1.11.2 // indirect github.com/mailru/easyjson v0.7.6 // indirect github.com/matttproud/golang_protobuf_extensions v1.0.2 // indirect + github.com/minio/crc64nvme v1.1.1 // indirect + github.com/minio/md5-simd v1.1.2 // indirect + github.com/minio/minio-go/v7 v7.0.99 // indirect github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect github.com/modern-go/reflect2 v1.0.2 // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect + github.com/philhofer/fwd v1.2.0 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/prometheus/client_golang v1.14.0 // indirect github.com/prometheus/client_model v0.3.0 // indirect github.com/prometheus/common v0.37.0 // indirect github.com/prometheus/procfs v0.8.0 // indirect + github.com/rs/xid v1.6.0 // indirect github.com/spf13/pflag v1.0.5 // indirect + github.com/tinylib/msgp v1.6.1 // indirect go.uber.org/atomic v1.7.0 // indirect go.uber.org/multierr v1.6.0 // indirect go.uber.org/zap v1.24.0 // indirect - golang.org/x/net v0.3.1-0.20221206200815-1e63c2f08a10 // indirect + go.yaml.in/yaml/v3 v3.0.4 // indirect + golang.org/x/crypto v0.46.0 // indirect + golang.org/x/net v0.48.0 // indirect golang.org/x/oauth2 v0.0.0-20220223155221-ee480838109b // indirect - golang.org/x/sys v0.3.0 // indirect - golang.org/x/term v0.3.0 // indirect - golang.org/x/text v0.5.0 // indirect + golang.org/x/sys v0.39.0 // indirect + golang.org/x/term v0.38.0 // indirect + golang.org/x/text v0.32.0 // indirect golang.org/x/time v0.3.0 // indirect gomodules.xyz/jsonpatch/v2 v2.2.0 // indirect google.golang.org/appengine v1.6.7 // indirect diff --git a/go.sum b/go.sum index 5e84df8..c599c83 100644 --- a/go.sum +++ b/go.sum @@ -58,6 +58,8 @@ github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSs github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/docopt/docopt-go v0.0.0-20180111231733-ee0de3bc6815/go.mod h1:WwZ+bS3ebgob9U8Nd0kOddGdZWjyMGR8Wziv+TBNwSE= +github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= +github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= github.com/emicklei/go-restful/v3 v3.9.0 h1:XwGDlfxEnQZzuopoqxwSEllNcCOM9DhhFyhFIIGKwxE= github.com/emicklei/go-restful/v3 v3.9.0/go.mod h1:6n3XBCmQQb25CM2LCACGz8ukIrRry+4bhvbpWn3mrbc= github.com/envoyproxy/go-control-plane v0.9.0/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4= @@ -73,6 +75,8 @@ github.com/fsnotify/fsnotify v1.6.0/go.mod h1:sl3t1tCWJFWoRz9R8WJCbQihKKwmorjAbS github.com/go-gl/glfw v0.0.0-20190409004039-e6da0acd62b1/go.mod h1:vR7hzQXu2zJy9AVAgeJqvqgH9Q5CA+iKCZ2gyEVpxRU= github.com/go-gl/glfw/v3.3/glfw v0.0.0-20191125211704-12ad95a8df72/go.mod h1:tQ2UAYgL5IevRw8kRxooKSPJfGvJ9fJQFa0TUsXzTg8= github.com/go-gl/glfw/v3.3/glfw v0.0.0-20200222043503-6f7a984d4dc4/go.mod h1:tQ2UAYgL5IevRw8kRxooKSPJfGvJ9fJQFa0TUsXzTg8= +github.com/go-ini/ini v1.67.0 h1:z6ZrTEZqSWOTyH2FlglNbNgARyHG8oLW9gMELqKr06A= +github.com/go-ini/ini v1.67.0/go.mod h1:ByCAeIL28uOIIG0E3PJtZPDL8WnHpFKFOtgjp+3Ies8= github.com/go-kit/kit v0.8.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as= github.com/go-kit/kit v0.9.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as= github.com/go-kit/log v0.1.0/go.mod h1:zbhenjAZHb184qTLMA9ZjW7ThYL0H2mk7Q6pNt4vbaY= @@ -159,8 +163,12 @@ github.com/google/pprof v0.0.0-20200708004538-1a94d8640e99/go.mod h1:ZgVRPoUq/hf github.com/google/renameio v0.1.0/go.mod h1:KWCgfxg9yswjAJkECMjeO8J8rahYeXnNhOm40UhjYkI= github.com/google/uuid v1.1.2 h1:EVhdT+1Kseyi1/pUmXKaFxYsDNy9RQYkMWRH68J/W7Y= github.com/google/uuid v1.1.2/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/googleapis/gax-go/v2 v2.0.4/go.mod h1:0Wqv26UfaUD9n4G6kQubkQ+KchISgw+vpHVxEJEs9eg= github.com/googleapis/gax-go/v2 v2.0.5/go.mod h1:DWXyrwAJ9X0FpwwEdw+IPEYBICEFu5mhpdKc/us6bOk= +github.com/gorilla/mux v1.8.1 h1:TuBL49tXwgrFYWhqrNgrUNEY92u81SPhu7sTdzQEiWY= +github.com/gorilla/mux v1.8.1/go.mod h1:AKf9I4AEqPTmMytcMc0KkNouC66V3BtZ4qD5fmWSiMQ= github.com/hashicorp/golang-lru v0.5.0/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8= github.com/hashicorp/golang-lru v0.5.1/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8= github.com/ianlancetaylor/demangle v0.0.0-20181102032728-5e5cf60278f6/go.mod h1:aSSvb/t6k1mPoxDqO4vJh6VOCGPwU4O0C2/Eqndh1Sc= @@ -181,6 +189,13 @@ github.com/julienschmidt/httprouter v1.2.0/go.mod h1:SYymIcj16QtmaHHD7aYtjjsJG7V github.com/julienschmidt/httprouter v1.3.0/go.mod h1:JR6WtHb+2LUe8TCKY3cZOxFyyO8IZAc4RVcycCCAKdM= github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8= github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= +github.com/klauspost/compress v1.18.2 h1:iiPHWW0YrcFgpBYhsA6D1+fqHssJscY/Tm/y2Uqnapk= +github.com/klauspost/compress v1.18.2/go.mod h1:R0h/fSBs8DE4ENlcrlib3PsXS61voFxhIs2DeRhCvJ4= +github.com/klauspost/cpuid/v2 v2.0.1/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg= +github.com/klauspost/cpuid/v2 v2.2.11 h1:0OwqZRYI2rFrjS4kvkDnqJkKHdHaRnCm68/DY4OxRzU= +github.com/klauspost/cpuid/v2 v2.2.11/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= +github.com/klauspost/crc32 v1.3.0 h1:sSmTt3gUt81RP655XGZPElI0PelVTZ6YwCRnPSupoFM= +github.com/klauspost/crc32 v1.3.0/go.mod h1:D7kQaZhnkX/Y0tstFGf8VUzv2UofNGqCjnC3zdHB0Hw= github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ= github.com/konsorten/go-windows-terminal-sequences v1.0.3/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ= github.com/kr/logfmt v0.0.0-20140226030751-b84e30acd515/go.mod h1:+0opPa2QZZtGFBFZlji/RkVcI2GknAs/DXo4wKdlNEc= @@ -190,6 +205,8 @@ github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/lib/pq v1.11.2 h1:x6gxUeu39V0BHZiugWe8LXZYZ+Utk7hSJGThs8sdzfs= +github.com/lib/pq v1.11.2/go.mod h1:/p+8NSbOcwzAEI7wiMXFlgydTwcgTr3OSKMsD2BitpA= github.com/mailru/easyjson v0.0.0-20190614124828-94de47d64c63/go.mod h1:C1wdFJiN94OJF2b5HbByQZoLdCWB1Yqtg26g4irojpc= github.com/mailru/easyjson v0.0.0-20190626092158-b2ccc519800e/go.mod h1:C1wdFJiN94OJF2b5HbByQZoLdCWB1Yqtg26g4irojpc= github.com/mailru/easyjson v0.7.6 h1:8yTIVnZgCoiM1TgqoeTl+LfU5Jg6/xL3QhGQnimLYnA= @@ -197,6 +214,12 @@ github.com/mailru/easyjson v0.7.6/go.mod h1:xzfreul335JAWq5oZzymOObrkdz5UnU4kGfJ github.com/matttproud/golang_protobuf_extensions v1.0.1/go.mod h1:D8He9yQNgCq6Z5Ld7szi9bcBfOoFv/3dc6xSMkL2PC0= github.com/matttproud/golang_protobuf_extensions v1.0.2 h1:hAHbPm5IJGijwng3PWk09JkG9WeqChjprR5s9bBZ+OM= github.com/matttproud/golang_protobuf_extensions v1.0.2/go.mod h1:BSXmuO+STAnVfrANrmjBb36TMTDstsz7MSK+HVaYKv4= +github.com/minio/crc64nvme v1.1.1 h1:8dwx/Pz49suywbO+auHCBpCtlW1OfpcLN7wYgVR6wAI= +github.com/minio/crc64nvme v1.1.1/go.mod h1:eVfm2fAzLlxMdUGc0EEBGSMmPwmXD5XiNRpnu9J3bvg= +github.com/minio/md5-simd v1.1.2 h1:Gdi1DZK69+ZVMoNHRXJyNcxrMA4dSxoYHZSQbirFg34= +github.com/minio/md5-simd v1.1.2/go.mod h1:MzdKDxYpY2BT9XQFocsiZf/NKVtR7nkE4RoEpN+20RM= +github.com/minio/minio-go/v7 v7.0.99 h1:2vH/byrwUkIpFQFOilvTfaUpvAX3fEFhEzO+DR3DlCE= +github.com/minio/minio-go/v7 v7.0.99/go.mod h1:EtGNKtlX20iL2yaYnxEigaIvj0G0GwSDnifnG8ClIdw= github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg= github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= @@ -214,6 +237,8 @@ github.com/onsi/ginkgo/v2 v2.6.0 h1:9t9b9vRUbFq3C4qKFCGkVuq/fIHji802N1nrtkh1mNc= github.com/onsi/ginkgo/v2 v2.6.0/go.mod h1:63DOGlLAH8+REH8jUGdL3YpCpu7JODesutUjdENfUAc= github.com/onsi/gomega v1.24.1 h1:KORJXNNTzJXzu4ScJWssJfJMnJ+2QJqhoQSRwNlze9E= github.com/onsi/gomega v1.24.1/go.mod h1:3AOiACssS3/MajrniINInwbfOOtfZvplPzuRSmvt1jM= +github.com/philhofer/fwd v1.2.0 h1:e6DnBTl7vGY+Gz322/ASL4Gyp1FspeMvx1RNDoToZuM= +github.com/philhofer/fwd v1.2.0/go.mod h1:RqIHx9QI14HlwKwm98g9Re5prTQ6LdeRQn+gXJFxsJM= github.com/pkg/errors v0.8.0/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= @@ -247,6 +272,8 @@ github.com/prometheus/procfs v0.7.3/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1 github.com/prometheus/procfs v0.8.0 h1:ODq8ZFEaYeCaZOJlZZdJA2AbQR98dSHSM1KW/You5mo= github.com/prometheus/procfs v0.8.0/go.mod h1:z7EfXMXOkbkqb9IINtpCn86r/to3BnA0uaxHdg830/4= github.com/rogpeppe/go-internal v1.3.0/go.mod h1:M8bDsm7K2OlrFYOpmOWEs/qY81heoFRclV5y23lUDJ4= +github.com/rs/xid v1.6.0 h1:fV591PaemRlL6JfRxGDEPl69wICngIQ3shQtzfy2gxU= +github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0= github.com/sirupsen/logrus v1.2.0/go.mod h1:LxeOpSwHxABJmUn/MG1IvRgCAasNZTLOkJPxbbu5VWo= github.com/sirupsen/logrus v1.4.2/go.mod h1:tLMulIdttU9McNUspp0xgXVQah82FyeX6MwdIuYE2rE= github.com/sirupsen/logrus v1.6.0/go.mod h1:7uNnSEd1DgxDLC74fIahvMZmmYsHGZGEOFrfsX/uA88= @@ -262,6 +289,8 @@ github.com/stretchr/testify v1.5.1/go.mod h1:5W2xD1RspED5o8YsWQXVCued0rvSQ+mT+I5 github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.8.0 h1:pSgiaMZlXftHpm5L7V1+rVB+AZJydKsMxsQBIJw4PKk= +github.com/tinylib/msgp v1.6.1 h1:ESRv8eL3u+DNHUoSAAQRE50Hm162zqAnBoGv9PzScPY= +github.com/tinylib/msgp v1.6.1/go.mod h1:RSp0LW9oSxFut3KzESt5Voq4GVWyS+PSulT77roAqEA= github.com/yuin/goldmark v1.1.25/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.1.32/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= @@ -280,12 +309,16 @@ go.uber.org/multierr v1.6.0/go.mod h1:cdWPpRnG4AhwMwsgIHip0KRBQjJy5kYEpYjJxpXp9i go.uber.org/zap v1.19.0/go.mod h1:xg/QME4nWcxGxrpdeYfq7UvYrLh66cuVKdrbD1XF/NI= go.uber.org/zap v1.24.0 h1:FiJd5l1UOLj0wCgbSE0rwwXHzEdAZS6hiiSnxJN/D60= go.uber.org/zap v1.24.0/go.mod h1:2kMP+WWQ8aoFoedH3T2sq6iJ2yDWpHbP0f6MQbS9Gkg= +go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= +go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= golang.org/x/crypto v0.0.0-20180904163835-0709b304e793/go.mod h1:6SG95UA2DQfeDnfUPMdvaQW0Q7yPrPDi9nlGo2tz2b4= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20190510104115-cbcb75029529/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20190605123033-f99c8df09eb5/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= +golang.org/x/crypto v0.46.0 h1:cKRW/pmt1pKAfetfu+RCEvjvZkA9RimPbh7bhFjGVBU= +golang.org/x/crypto v0.46.0/go.mod h1:Evb/oLKmMraqjZ2iQTwDwvCtJkczlDuTmdJXoZVzqU0= golang.org/x/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA= golang.org/x/exp v0.0.0-20190306152737-a1d7652674e8/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA= golang.org/x/exp v0.0.0-20190510132918-efd6b22b2522/go.mod h1:ZjyILWgesfNpC6sMxTJOJm9Kp84zZh5NQWvqDGG3Qr8= @@ -350,6 +383,8 @@ golang.org/x/net v0.0.0-20220127200216-cd36cc0744dd/go.mod h1:CfG3xpIq0wQ8r1q4Su golang.org/x/net v0.0.0-20220225172249-27dd8689420f/go.mod h1:CfG3xpIq0wQ8r1q4Su4UZFWDARRcnwPjda9FqA0JpMk= golang.org/x/net v0.3.1-0.20221206200815-1e63c2f08a10 h1:Frnccbp+ok2GkUS2tC84yAq/U9Vg+0sIO7aRL3T4Xnc= golang.org/x/net v0.3.1-0.20221206200815-1e63c2f08a10/go.mod h1:MBQ8lrhLObU/6UmLb4fmbmk5OcyYmqtbGd/9yIeKjEE= +golang.org/x/net v0.48.0 h1:zyQRTTrjc33Lhh0fBgT/H3oZq9WuvRR5gPC70xpDiQU= +golang.org/x/net v0.48.0/go.mod h1:+ndRgGjkh8FGtu1w1FGbEC31if4VrNVMuKTgcAAnQRY= golang.org/x/oauth2 v0.0.0-20180821212333-d2e6202438be/go.mod h1:N/0e6XlmueqKjAGxoOufVs8QHGRruUQn6yWY3a++T0U= golang.org/x/oauth2 v0.0.0-20190226205417-e64efc72b421/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw= golang.org/x/oauth2 v0.0.0-20190604053449-0f29369cfe45/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw= @@ -410,10 +445,14 @@ golang.org/x/sys v0.0.0-20220114195835-da31bd327af9/go.mod h1:oPkhp1MJrh7nUepCBc golang.org/x/sys v0.0.0-20220908164124-27713097b956/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.3.0 h1:w8ZOecv6NaNa/zC8944JTU3vz4u6Lagfk4RPQxv92NQ= golang.org/x/sys v0.3.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.39.0 h1:CvCKL8MeisomCi6qNZ+wbb0DN9E5AATixKsvNtMoMFk= +golang.org/x/sys v0.39.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= golang.org/x/term v0.3.0 h1:qoo4akIqOcDME5bhc/NgxUdovd6BSS2uMsVjB56q1xI= golang.org/x/term v0.3.0/go.mod h1:q750SLmJuPmVoN1blW3UFBPREJfb1KmY3vwxfr+nFDA= +golang.org/x/term v0.38.0 h1:PQ5pkm/rLO6HnxFR7N2lJHOZX6Kez5Y1gDSJla6jo7Q= +golang.org/x/term v0.38.0/go.mod h1:bSEAKrOT1W+VSu9TSCMtoGEOUcKxOKgl3LE5QEF/xVg= golang.org/x/text v0.0.0-20170915032832-14c0d48ead0c/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.1-0.20180807135948-17ff2d5776d2/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= @@ -423,6 +462,8 @@ golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= golang.org/x/text v0.5.0 h1:OLmvp0KP+FVG99Ct/qFiL/Fhk4zp4QQnZ7b2U+5piUM= golang.org/x/text v0.5.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= +golang.org/x/text v0.32.0 h1:ZD01bjUt1FQ9WJ0ClOL5vxgxOI/sVCNgX1YtKwcY0mU= +golang.org/x/text v0.32.0/go.mod h1:o/rUWzghvpD5TXrTIBuJU77MTaN0ljMWE47kxGJQ7jY= golang.org/x/time v0.0.0-20181108054448-85acf8d2951c/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= golang.org/x/time v0.0.0-20190308202827-9d24e82272b4/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= golang.org/x/time v0.0.0-20191024005414-555d28b269f0/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= diff --git a/internal/api/handler/functions.go b/internal/api/handler/functions.go new file mode 100644 index 0000000..da35055 --- /dev/null +++ b/internal/api/handler/functions.go @@ -0,0 +1,193 @@ +// Изменено: 2026-03-07 +// functions.go — CRUD handlers для Function CRD. +// Принимает JSON, создаёт/обновляет/удаляет k8s ресурсы Function. +// Namespace берётся из URL: /v1/namespaces/{namespace}/functions/{name} + +package handler + +import ( + "encoding/json" + "net/http" + + "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" + + slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1" +) + +// functionRequest — тело запроса для создания/обновления функции. +type functionRequest struct { + Name string `json:"name"` + Runtime string `json:"runtime"` + Entrypoint string `json:"entrypoint"` + MemoryMB int32 `json:"memory_mb"` + TimeoutSec int32 `json:"timeout_sec"` + Env map[string]string `json:"env_vars"` + S3Bucket string `json:"s3_bucket"` + S3Key string `json:"s3_key"` +} + +// functionResponse — ответ при чтении функции. +type functionResponse struct { + Name string `json:"name"` + Namespace string `json:"namespace"` + Runtime string `json:"runtime"` + Entrypoint string `json:"entrypoint"` + MemoryMB int32 `json:"memory_mb"` + TimeoutSec int32 `json:"timeout_sec"` + Env map[string]string `json:"env_vars"` + S3Bucket string `json:"s3_bucket"` + S3Key string `json:"s3_key"` + Phase slessv1alpha1.FunctionPhase `json:"phase"` + ImageRef string `json:"image_ref"` + Message string `json:"message,omitempty"` +} + +// fnToResponse конвертирует CRD в ответ API. +func fnToResponse(fn *slessv1alpha1.Function) functionResponse { + return functionResponse{ + Name: fn.Name, + Namespace: fn.Namespace, + Runtime: fn.Spec.Runtime, + Entrypoint: fn.Spec.Entrypoint, + MemoryMB: fn.Spec.MemoryMB, + TimeoutSec: fn.Spec.TimeoutSec, + Env: fn.Spec.Env, + S3Bucket: fn.Spec.S3Bucket, + S3Key: fn.Spec.S3Key, + Phase: fn.Status.Phase, + ImageRef: fn.Status.ImageRef, + Message: fn.Status.Message, + } +} + +// ListFunctions — GET /v1/namespaces/{namespace}/functions +func (h *Handler) ListFunctions(w http.ResponseWriter, r *http.Request) { + ns := namespace(r) + list := &slessv1alpha1.FunctionList{} + if err := h.K8s.List(r.Context(), list, client.InNamespace(ns)); err != nil { + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + result := make([]functionResponse, 0, len(list.Items)) + for i := range list.Items { + result = append(result, fnToResponse(&list.Items[i])) + } + writeJSON(w, http.StatusOK, result) +} + +// CreateFunction — POST /v1/namespaces/{namespace}/functions +func (h *Handler) CreateFunction(w http.ResponseWriter, r *http.Request) { + ns := namespace(r) + var req functionRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeJSON(w, http.StatusBadRequest, errResp("invalid JSON: "+err.Error())) + return + } + if req.Name == "" || req.Runtime == "" { + writeJSON(w, http.StatusBadRequest, errResp("name and runtime are required")) + return + } + + fn := &slessv1alpha1.Function{ + ObjectMeta: metav1.ObjectMeta{ + Name: req.Name, + Namespace: ns, + }, + Spec: slessv1alpha1.FunctionSpec{ + Runtime: req.Runtime, + Entrypoint: req.Entrypoint, + MemoryMB: req.MemoryMB, + TimeoutSec: req.TimeoutSec, + Env: req.Env, + S3Bucket: req.S3Bucket, + S3Key: req.S3Key, + }, + } + if err := h.K8s.Create(r.Context(), fn); err != nil { + if errors.IsAlreadyExists(err) { + writeJSON(w, http.StatusConflict, errResp("function already exists")) + return + } + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + writeJSON(w, http.StatusCreated, fnToResponse(fn)) +} + +// GetFunction — GET /v1/namespaces/{namespace}/functions/{name} +func (h *Handler) GetFunction(w http.ResponseWriter, r *http.Request) { + ns := namespace(r) + name := pathVar(r, "name") + fn := &slessv1alpha1.Function{} + if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, fn); err != nil { + if errors.IsNotFound(err) { + writeJSON(w, http.StatusNotFound, errResp("function not found")) + return + } + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + writeJSON(w, http.StatusOK, fnToResponse(fn)) +} + +// UpdateFunction — PUT /v1/namespaces/{namespace}/functions/{name} +func (h *Handler) UpdateFunction(w http.ResponseWriter, r *http.Request) { + ns := namespace(r) + name := pathVar(r, "name") + var req functionRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeJSON(w, http.StatusBadRequest, errResp("invalid JSON: "+err.Error())) + return + } + + fn := &slessv1alpha1.Function{} + if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, fn); err != nil { + if errors.IsNotFound(err) { + writeJSON(w, http.StatusNotFound, errResp("function not found")) + return + } + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + + // Обновляем только изменяемые поля spec + fn.Spec.Runtime = req.Runtime + fn.Spec.Entrypoint = req.Entrypoint + fn.Spec.MemoryMB = req.MemoryMB + fn.Spec.TimeoutSec = req.TimeoutSec + fn.Spec.Env = req.Env + if req.S3Bucket != "" { + fn.Spec.S3Bucket = req.S3Bucket + } + if req.S3Key != "" { + fn.Spec.S3Key = req.S3Key + } + + if err := h.K8s.Update(r.Context(), fn); err != nil { + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + writeJSON(w, http.StatusOK, fnToResponse(fn)) +} + +// DeleteFunction — DELETE /v1/namespaces/{namespace}/functions/{name} +func (h *Handler) DeleteFunction(w http.ResponseWriter, r *http.Request) { + ns := namespace(r) + name := pathVar(r, "name") + fn := &slessv1alpha1.Function{} + if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, fn); err != nil { + if errors.IsNotFound(err) { + w.WriteHeader(http.StatusNoContent) + return + } + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + if err := h.K8s.Delete(r.Context(), fn); err != nil { + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + w.WriteHeader(http.StatusNoContent) +} diff --git a/internal/api/handler/handler.go b/internal/api/handler/handler.go new file mode 100644 index 0000000..c7c51c5 --- /dev/null +++ b/internal/api/handler/handler.go @@ -0,0 +1,53 @@ +// Изменено: 2026-03-07 +// Handler — общий контейнер зависимостей для всех REST handlers. +// Все handlers получают доступ к k8s, S3 и Postgres через эту структуру. +// Логирование через slog, маршрутизация через gorilla/mux. + +package handler + +import ( + "encoding/json" + "log/slog" + "net/http" + + "github.com/gorilla/mux" + "k8s.io/apimachinery/pkg/runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + + "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/storage/postgres" + "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/storage/s3" +) + +// Handler содержит зависимости для всех REST-обработчиков. +type Handler struct { + K8s client.Client + Scheme *runtime.Scheme + S3 *s3.Client + PG *postgres.Store + Log *slog.Logger +} + +// writeJSON отправляет JSON-ответ с указанным статусом. +func writeJSON(w http.ResponseWriter, status int, v any) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(v) +} + +// errResp возвращает структуру ошибки для JSON. +func errResp(msg string) map[string]string { + return map[string]string{"error": msg} +} + +// pathVar читает переменную пути из gorilla/mux. +func pathVar(r *http.Request, key string) string { + return mux.Vars(r)[key] +} + +// namespace читает {namespace} из пути, fallback — "default". +func namespace(r *http.Request) string { + if ns := mux.Vars(r)["namespace"]; ns != "" { + return ns + } + return "default" +} diff --git a/internal/api/handler/invocations.go b/internal/api/handler/invocations.go new file mode 100644 index 0000000..369b7ac --- /dev/null +++ b/internal/api/handler/invocations.go @@ -0,0 +1,61 @@ +// Изменено: 2026-03-07 +// invocations.go — GET-handlers для логов вызовов функций. +// Данные читаются из PostgreSQL (invocations таблица). +// POST /invoke остаётся для будущего прямого вызова (v2). + +package handler + +import ( + "net/http" + "strconv" +) + +// invocationResponse — ответ при чтении записи вызова. +type invocationResponse struct { + ID string `json:"id"` + FunctionName string `json:"function_name"` + Namespace string `json:"namespace"` + Status string `json:"status"` + DurationMs int32 `json:"duration_ms"` + HTTPStatus *int32 `json:"http_status,omitempty"` // nil для cron триггеров + Logs string `json:"logs,omitempty"` + TriggerType string `json:"trigger_type"` + CreatedAt string `json:"created_at"` +} + +// ListInvocations — GET /v1/namespaces/{namespace}/functions/{name}/invocations +// Возвращает последние limit вызовов (default 50, max 200). +func (h *Handler) ListInvocations(w http.ResponseWriter, r *http.Request) { + ns := namespace(r) + fnName := pathVar(r, "name") + + limit := 50 + if q := r.URL.Query().Get("limit"); q != "" { + if v, err := strconv.Atoi(q); err == nil && v > 0 && v <= 200 { + limit = v + } + } + + rows, err := h.PG.ListInvocations(r.Context(), fnName, ns, limit) + if err != nil { + h.Log.Error("list invocations", "fn", fnName, "ns", ns, "err", err) + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + + result := make([]invocationResponse, 0, len(rows)) + for _, row := range rows { + result = append(result, invocationResponse{ + ID: row.ID, + FunctionName: row.FunctionName, + Namespace: row.Namespace, + Status: row.Status, + DurationMs: row.DurationMs, + HTTPStatus: row.HTTPStatus, + Logs: row.Logs, + TriggerType: row.TriggerType, + CreatedAt: row.CreatedAt.Format("2006-01-02T15:04:05Z"), + }) + } + writeJSON(w, http.StatusOK, result) +} diff --git a/internal/api/handler/triggers.go b/internal/api/handler/triggers.go new file mode 100644 index 0000000..3807e48 --- /dev/null +++ b/internal/api/handler/triggers.go @@ -0,0 +1,143 @@ +// Изменено: 2026-03-07 +// triggers.go — CRUD handlers для Trigger CRD. +// Триггеры привязаны к Function через FunctionRef. +// Namespace берётся из URL: /v1/namespaces/{namespace}/triggers/{name} + +package handler + +import ( + "encoding/json" + "net/http" + + "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" + + slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1" +) + +// triggerRequest — тело запроса для создания триггера. +type triggerRequest struct { + Name string `json:"name"` + Type string `json:"type"` // http | cron + FunctionRef string `json:"function"` // имя Function CRD + Schedule string `json:"schedule"` // cron-расписание, только для type=cron + PreWarmSeconds int32 `json:"pre_warm_seconds"` +} + +// triggerResponse — ответ при чтении триггера. +type triggerResponse struct { + Name string `json:"name"` + Namespace string `json:"namespace"` + Type string `json:"type"` + FunctionRef string `json:"function"` + Schedule string `json:"schedule,omitempty"` + Active bool `json:"active"` + URL string `json:"url,omitempty"` + Message string `json:"message,omitempty"` +} + +// trToResponse конвертирует Trigger CRD в ответ API. +func trToResponse(tr *slessv1alpha1.Trigger) triggerResponse { + return triggerResponse{ + Name: tr.Name, + Namespace: tr.Namespace, + Type: string(tr.Spec.Type), + FunctionRef: tr.Spec.FunctionRef, + Schedule: tr.Spec.Schedule, + Active: tr.Status.Active, + URL: tr.Status.URL, + Message: tr.Status.Message, + } +} + +// ListTriggers — GET /v1/namespaces/{namespace}/triggers +func (h *Handler) ListTriggers(w http.ResponseWriter, r *http.Request) { + ns := namespace(r) + list := &slessv1alpha1.TriggerList{} + if err := h.K8s.List(r.Context(), list, client.InNamespace(ns)); err != nil { + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + result := make([]triggerResponse, 0, len(list.Items)) + for i := range list.Items { + result = append(result, trToResponse(&list.Items[i])) + } + writeJSON(w, http.StatusOK, result) +} + +// CreateTrigger — POST /v1/namespaces/{namespace}/triggers +func (h *Handler) CreateTrigger(w http.ResponseWriter, r *http.Request) { + ns := namespace(r) + var req triggerRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeJSON(w, http.StatusBadRequest, errResp("invalid JSON: "+err.Error())) + return + } + if req.Name == "" || req.Type == "" || req.FunctionRef == "" { + writeJSON(w, http.StatusBadRequest, errResp("name, type and function are required")) + return + } + if req.Type == "cron" && req.Schedule == "" { + writeJSON(w, http.StatusBadRequest, errResp("schedule is required for cron trigger")) + return + } + + tr := &slessv1alpha1.Trigger{ + ObjectMeta: metav1.ObjectMeta{ + Name: req.Name, + Namespace: ns, + }, + Spec: slessv1alpha1.TriggerSpec{ + Type: slessv1alpha1.TriggerType(req.Type), + FunctionRef: req.FunctionRef, + Schedule: req.Schedule, + PreWarmSeconds: req.PreWarmSeconds, + }, + } + if err := h.K8s.Create(r.Context(), tr); err != nil { + if errors.IsAlreadyExists(err) { + writeJSON(w, http.StatusConflict, errResp("trigger already exists")) + return + } + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + writeJSON(w, http.StatusCreated, trToResponse(tr)) +} + +// GetTrigger — GET /v1/namespaces/{namespace}/triggers/{name} +func (h *Handler) GetTrigger(w http.ResponseWriter, r *http.Request) { + ns := namespace(r) + name := pathVar(r, "name") + tr := &slessv1alpha1.Trigger{} + if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, tr); err != nil { + if errors.IsNotFound(err) { + writeJSON(w, http.StatusNotFound, errResp("trigger not found")) + return + } + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + writeJSON(w, http.StatusOK, trToResponse(tr)) +} + +// DeleteTrigger — DELETE /v1/namespaces/{namespace}/triggers/{name} +func (h *Handler) DeleteTrigger(w http.ResponseWriter, r *http.Request) { + ns := namespace(r) + name := pathVar(r, "name") + tr := &slessv1alpha1.Trigger{} + if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, tr); err != nil { + if errors.IsNotFound(err) { + w.WriteHeader(http.StatusNoContent) + return + } + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + if err := h.K8s.Delete(r.Context(), tr); err != nil { + writeJSON(w, http.StatusInternalServerError, errResp(err.Error())) + return + } + w.WriteHeader(http.StatusNoContent) +} diff --git a/internal/api/middleware/auth.go b/internal/api/middleware/auth.go new file mode 100644 index 0000000..6b99ecd --- /dev/null +++ b/internal/api/middleware/auth.go @@ -0,0 +1,40 @@ +// Изменено: 2026-03-07 +// Auth middleware — проверяет Bearer токен из заголовка Authorization. +// В v1: сравниваем с SLESS_API_TOKEN из конфига. +// В prod: вызываем auth-сервис nubes.ru (TODO v2). + +package middleware + +import ( + "log/slog" + "net/http" + "strings" +) + +// Auth возвращает middleware которое требует заголовок: +// +// Authorization: Bearer +// +// и сравнивает его с allowedToken. +func Auth(allowedToken string, log *slog.Logger, next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + header := r.Header.Get("Authorization") + if header == "" { + log.Warn("auth: missing authorization header", "remote", r.RemoteAddr, "path", r.URL.Path) + http.Error(w, `{"error":"authorization required"}`, http.StatusUnauthorized) + return + } + parts := strings.SplitN(header, " ", 2) + if len(parts) != 2 || !strings.EqualFold(parts[0], "bearer") { + log.Warn("auth: invalid authorization format", "remote", r.RemoteAddr) + http.Error(w, `{"error":"invalid authorization format, use Bearer "}`, http.StatusUnauthorized) + return + } + if parts[1] != allowedToken { + log.Warn("auth: invalid token", "remote", r.RemoteAddr, "path", r.URL.Path) + http.Error(w, `{"error":"invalid token"}`, http.StatusForbidden) + return + } + next.ServeHTTP(w, r) + }) +} diff --git a/internal/api/middleware/logging.go b/internal/api/middleware/logging.go new file mode 100644 index 0000000..184c6d5 --- /dev/null +++ b/internal/api/middleware/logging.go @@ -0,0 +1,38 @@ +// Изменено: 2026-03-07 +// Logging middleware — логирует каждый HTTP запрос через slog. +// Записывает метод, путь, статус и длительность. + +package middleware + +import ( + "log/slog" + "net/http" + "time" +) + +// responseWriter оборачивает http.ResponseWriter чтобы захватить статус ответа. +type responseWriter struct { + http.ResponseWriter + status int +} + +func (rw *responseWriter) WriteHeader(code int) { + rw.status = code + rw.ResponseWriter.WriteHeader(code) +} + +// Logging возвращает middleware которое логирует все запросы через slog. +func Logging(log *slog.Logger, next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + start := time.Now() + rw := &responseWriter{ResponseWriter: w, status: http.StatusOK} + next.ServeHTTP(rw, r) + log.Info("http", + "method", r.Method, + "path", r.URL.Path, + "status", rw.status, + "duration_ms", time.Since(start).Milliseconds(), + "remote", r.RemoteAddr, + ) + }) +} diff --git a/internal/api/router.go b/internal/api/router.go new file mode 100644 index 0000000..28b6619 --- /dev/null +++ b/internal/api/router.go @@ -0,0 +1,48 @@ +// Изменено: 2026-03-07 +// router.go — регистрация всех REST-маршрутов через gorilla/mux. +// Все маршруты защищены Bearer-токеном (middleware.Auth). +// Маршруты сгруппированы по /v1/namespaces/{namespace}/... + +package api + +import ( + "log/slog" + "net/http" + + "github.com/gorilla/mux" + + "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/api/handler" + "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/api/middleware" +) + +// NewRouter собирает gorilla/mux роутер со всеми маршрутами. +// apiToken — статический Bearer-токен для v1 аутентификации. +func NewRouter(h *handler.Handler, apiToken string, log *slog.Logger) http.Handler { + r := mux.NewRouter() + + // Суброутер для /v1 — все маршруты API + v1 := r.PathPrefix("/v1").Subrouter() + + // Functions CRUD + v1.HandleFunc("/namespaces/{namespace}/functions", h.ListFunctions).Methods(http.MethodGet) + v1.HandleFunc("/namespaces/{namespace}/functions", h.CreateFunction).Methods(http.MethodPost) + v1.HandleFunc("/namespaces/{namespace}/functions/{name}", h.GetFunction).Methods(http.MethodGet) + v1.HandleFunc("/namespaces/{namespace}/functions/{name}", h.UpdateFunction).Methods(http.MethodPut) + v1.HandleFunc("/namespaces/{namespace}/functions/{name}", h.DeleteFunction).Methods(http.MethodDelete) + + // Invocation logs + v1.HandleFunc("/namespaces/{namespace}/functions/{name}/invocations", h.ListInvocations).Methods(http.MethodGet) + + // Triggers CRUD + v1.HandleFunc("/namespaces/{namespace}/triggers", h.ListTriggers).Methods(http.MethodGet) + v1.HandleFunc("/namespaces/{namespace}/triggers", h.CreateTrigger).Methods(http.MethodPost) + v1.HandleFunc("/namespaces/{namespace}/triggers/{name}", h.GetTrigger).Methods(http.MethodGet) + v1.HandleFunc("/namespaces/{namespace}/triggers/{name}", h.DeleteTrigger).Methods(http.MethodDelete) + + // Цепочка middleware: logging → auth → router + // Порядок важен: сначала логируем (чтобы видеть все запросы включая отклонённые), + // затем проверяем авторизацию. + return middleware.Logging(log, + middleware.Auth(apiToken, log, r), + ) +} diff --git a/internal/builder/.gitkeep b/internal/builder/.gitkeep deleted file mode 100644 index e69de29..0000000 diff --git a/internal/builder/builder.go b/internal/builder/builder.go new file mode 100644 index 0000000..d897d43 --- /dev/null +++ b/internal/builder/builder.go @@ -0,0 +1,155 @@ +// Изменено: 2026-03-07 +// Builder — запускает kaniko Job в k8s для сборки Docker образа из кода функции. +// Почему kaniko, а не docker-in-docker (DinD): +// kaniko не требует privileged контейнер, что безопаснее в managed кластере. +// kaniko читает контекст сборки из S3 напрямую. +// Workflow: FunctionController вызывает Build → Job запускается → образ пушится в registry. + +package builder + +import ( + "context" + "fmt" + "time" + + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// Builder — управляет сборкой Docker образов через kaniko Jobs в k8s. +type Builder struct { + client client.Client + builderImage string // образ kaniko + registryHost string // куда пушим образ + s3Endpoint string // откуда kaniko берёт код + s3AccessKey string + s3SecretKey string + s3Bucket string + namespace string // namespace где запускаем build Job'ы +} + +// Config — параметры для создания Builder'а. +type Config struct { + BuilderImage string + RegistryHost string + S3Endpoint string + S3AccessKey string + S3SecretKey string + S3Bucket string + Namespace string +} + +// New создаёт новый Builder. +func New(c client.Client, cfg Config) *Builder { + return &Builder{ + client: c, + builderImage: cfg.BuilderImage, + registryHost: cfg.RegistryHost, + s3Endpoint: cfg.S3Endpoint, + s3AccessKey: cfg.S3AccessKey, + s3SecretKey: cfg.S3SecretKey, + s3Bucket: cfg.S3Bucket, + namespace: cfg.Namespace, + } +} + +// ImageRef возвращает полный путь к образу в registry для данной функции и версии. +func (b *Builder) ImageRef(namespace, funcName, s3Key string) string { + // Используем s3Key как уникальный тег чтобы разные версии не перезаписывали друг друга + return fmt.Sprintf("%s/sless-%s-%s:latest", b.registryHost, namespace, funcName) +} + +// Build запускает kaniko Job для сборки образа функции. +// Возвращает имя Job'а чтобы контроллер мог следить за его статусом. +func (b *Builder) Build(ctx context.Context, namespace, funcName, s3Key string) (string, error) { + imageRef := b.ImageRef(namespace, funcName, s3Key) + jobName := fmt.Sprintf("build-%s-%s", funcName, time.Now().Format("20060102150405")) + + // kaniko читает Dockerfile из context архива в S3 + // --context=s3://bucket/key — kaniko поддерживает S3 как источник контекста + s3ContextURL := fmt.Sprintf("s3://%s/%s", b.s3Bucket, s3Key) + + job := &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{ + Name: jobName, + Namespace: b.namespace, + Labels: map[string]string{ + "app": "sless-builder", + "function-name": funcName, + "function-ns": namespace, + }, + }, + Spec: batchv1.JobSpec{ + // Не повторяем при ошибке — контроллер сам перезапустит reconcile + BackoffLimit: int32Ptr(0), + Completions: int32Ptr(1), + Template: corev1.PodTemplateSpec{ + Spec: corev1.PodSpec{ + RestartPolicy: corev1.RestartPolicyNever, + Containers: []corev1.Container{ + { + Name: "kaniko", + Image: b.builderImage, + Args: []string{ + "--context=" + s3ContextURL, + "--destination=" + imageRef, + "--skip-tls-verify", // для внутреннего registry без TLS + "--cache=true", // кешируем слои для ускорения повторных сборок + }, + Env: []corev1.EnvVar{ + {Name: "AWS_ACCESS_KEY_ID", Value: b.s3AccessKey}, + {Name: "AWS_SECRET_ACCESS_KEY", Value: b.s3SecretKey}, + {Name: "S3_ENDPOINT", Value: b.s3Endpoint}, + // Kaniko использует AWS SDK совместимый с S3 — задаём кастомный endpoint + {Name: "AWS_REGION", Value: "us-east-1"}, + }, + }, + }, + }, + }, + }, + } + + if err := b.client.Create(ctx, job); err != nil { + return "", fmt.Errorf("create build job: %w", err) + } + return jobName, nil +} + +// JobStatus проверяет статус build Job'а. +// Возвращает: "running", "succeeded", "failed" +func (b *Builder) JobStatus(ctx context.Context, jobName string) (string, error) { + job := &batchv1.Job{} + if err := b.client.Get(ctx, client.ObjectKey{Name: jobName, Namespace: b.namespace}, job); err != nil { + if errors.IsNotFound(err) { + return "failed", nil + } + return "", fmt.Errorf("get build job: %w", err) + } + if job.Status.Succeeded > 0 { + return "succeeded", nil + } + if job.Status.Failed > 0 { + return "failed", nil + } + return "running", nil +} + +// Cleanup удаляет завершённый build Job из k8s. +// Вызывается после того как контроллер зафиксировал результат сборки. +func (b *Builder) Cleanup(ctx context.Context, jobName string) error { + job := &batchv1.Job{} + if err := b.client.Get(ctx, client.ObjectKey{Name: jobName, Namespace: b.namespace}, job); err != nil { + if errors.IsNotFound(err) { + return nil // уже удалён + } + return fmt.Errorf("get build job for cleanup: %w", err) + } + propagation := metav1.DeletePropagationForeground + return b.client.Delete(ctx, job, &client.DeleteOptions{PropagationPolicy: &propagation}) +} + +func int32Ptr(i int32) *int32 { return &i } diff --git a/internal/config/.gitkeep b/internal/config/.gitkeep deleted file mode 100644 index e69de29..0000000 diff --git a/internal/config/config.go b/internal/config/config.go new file mode 100644 index 0000000..3529ebb --- /dev/null +++ b/internal/config/config.go @@ -0,0 +1,122 @@ +// Изменено: 2026-03-07 +// Конфигурация сервиса — читается из env переменных при старте. +// Все компоненты (API, builder, runner) получают конфиг через эту структуру. +// Используем env а не файлы конфигурации — стандарт для k8s (ConfigMap/Secret → env). + +package config + +import ( + "fmt" + "os" + "strconv" +) + +// Config — централизованная конфигурация всего сервиса. +type Config struct { + // APIPort — порт на котором слушает REST API сервер (default: 8080) + APIPort int + + // PostgreSQL — параметры подключения к БД для хранения логов вызовов + PostgresDSN string + + // S3 — параметры подключения к Ceph S3 для хранения кода функций + S3Endpoint string + S3AccessKey string + S3SecretKey string + S3Bucket string + S3UseSSL bool + + // Registry — адрес внутреннего Docker registry куда пушим собранные образы + RegistryHost string + + // FunctionNamespace — префикс namespace для функций пользователей. + // Итоговый namespace: FunctionNamespacePrefix + "-" + functionName + FunctionNamespacePrefix string + + // BuilderImage — образ для сборки функций (kaniko или buildah) + BuilderImage string + + // IngressHost — базовый домен для HTTP триггеров, например: fn.kube5s.ru + // URL функции будет: https://{funcName}-{namespace}.{IngressHost} + IngressHost string + + // APIToken — статический Bearer-токен для v1 REST API аутентификации. + // В prod заменить на вызов auth-сервиса. + APIToken string +} + +// Load читает конфиг из env переменных. +// Возвращает ошибку если обязательные переменные не заданы. +func Load() (*Config, error) { + cfg := &Config{ + APIPort: 8080, + S3UseSSL: false, + FunctionNamespacePrefix: "sless-fn", + BuilderImage: "gcr.io/kaniko-project/executor:latest", + } + + // Опциональный порт API + if v := os.Getenv("API_PORT"); v != "" { + port, err := strconv.Atoi(v) + if err != nil { + return nil, fmt.Errorf("API_PORT must be a number: %w", err) + } + cfg.APIPort = port + } + + // Обязательные параметры PostgreSQL + cfg.PostgresDSN = os.Getenv("POSTGRES_DSN") + if cfg.PostgresDSN == "" { + return nil, fmt.Errorf("POSTGRES_DSN is required") + } + + // Обязательные параметры S3 + cfg.S3Endpoint = os.Getenv("S3_ENDPOINT") + if cfg.S3Endpoint == "" { + return nil, fmt.Errorf("S3_ENDPOINT is required") + } + cfg.S3AccessKey = os.Getenv("S3_ACCESS_KEY") + if cfg.S3AccessKey == "" { + return nil, fmt.Errorf("S3_ACCESS_KEY is required") + } + cfg.S3SecretKey = os.Getenv("S3_SECRET_KEY") + if cfg.S3SecretKey == "" { + return nil, fmt.Errorf("S3_SECRET_KEY is required") + } + cfg.S3Bucket = os.Getenv("S3_BUCKET") + if cfg.S3Bucket == "" { + // Используем дефолтный бакет если не задан + cfg.S3Bucket = "sless-functions" + } + if v := os.Getenv("S3_USE_SSL"); v == "true" { + cfg.S3UseSSL = true + } + + // Обязательный адрес registry + cfg.RegistryHost = os.Getenv("REGISTRY_HOST") + if cfg.RegistryHost == "" { + return nil, fmt.Errorf("REGISTRY_HOST is required") + } + + // Опциональные параметры с дефолтами + if v := os.Getenv("FUNCTION_NAMESPACE_PREFIX"); v != "" { + cfg.FunctionNamespacePrefix = v + } + if v := os.Getenv("BUILDER_IMAGE"); v != "" { + cfg.BuilderImage = v + } + + // IngressHost — базовый домен для HTTP триггеров + cfg.IngressHost = os.Getenv("INGRESS_HOST") + if cfg.IngressHost == "" { + cfg.IngressHost = "fn.kube5s.ru" + } + + // APIToken — обязательный токен для REST API + cfg.APIToken = os.Getenv("SLESS_API_TOKEN") + if cfg.APIToken == "" { + return nil, fmt.Errorf("SLESS_API_TOKEN is required") + } + + return cfg, nil +} diff --git a/internal/storage/postgres/.gitkeep b/internal/storage/postgres/.gitkeep deleted file mode 100644 index e69de29..0000000 diff --git a/internal/storage/postgres/store.go b/internal/storage/postgres/store.go new file mode 100644 index 0000000..279fbe6 --- /dev/null +++ b/internal/storage/postgres/store.go @@ -0,0 +1,108 @@ +// Изменено: 2026-03-07 +// PostgreSQL storage — хранение логов вызовов функций. +// Состояние функций (phase, imageRef) хранится в k8s CRD, не здесь. +// Здесь только invocations: логи, статусы, время выполнения. + +package postgres + +import ( + "context" + "database/sql" + "fmt" + "time" + + _ "github.com/lib/pq" // драйвер PostgreSQL +) + +// Invocation — запись о вызове функции. +type Invocation struct { + ID string + FunctionName string + Namespace string + Status string // success, error, timeout + DurationMs int32 + HTTPStatus *int32 // nil для cron триггеров + Logs string + TriggerType string + CreatedAt time.Time +} + +// Store — клиент для работы с PostgreSQL. +type Store struct { + db *sql.DB +} + +// New открывает подключение к PostgreSQL и проверяет его. +func New(dsn string) (*Store, error) { + db, err := sql.Open("postgres", dsn) + if err != nil { + return nil, fmt.Errorf("open postgres: %w", err) + } + // Проверяем что подключение реально работает + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := db.PingContext(ctx); err != nil { + return nil, fmt.Errorf("ping postgres: %w", err) + } + db.SetMaxOpenConns(10) + db.SetMaxIdleConns(5) + db.SetConnMaxLifetime(5 * time.Minute) + return &Store{db: db}, nil +} + +// Close закрывает подключение к БД. +func (s *Store) Close() error { + return s.db.Close() +} + +// SaveInvocation записывает лог вызова функции в БД. +func (s *Store) SaveInvocation(ctx context.Context, inv *Invocation) error { + _, err := s.db.ExecContext(ctx, ` + INSERT INTO invocations (function_name, namespace, status, duration_ms, http_status, logs, trigger_type) + VALUES ($1, $2, $3, $4, $5, $6, $7) + `, inv.FunctionName, inv.Namespace, inv.Status, inv.DurationMs, inv.HTTPStatus, inv.Logs, inv.TriggerType) + if err != nil { + return fmt.Errorf("save invocation: %w", err) + } + return nil +} + +// ListInvocations возвращает последние N вызовов для указанной функции. +func (s *Store) ListInvocations(ctx context.Context, functionName, namespace string, limit int) ([]*Invocation, error) { + rows, err := s.db.QueryContext(ctx, ` + SELECT id, function_name, namespace, status, duration_ms, http_status, logs, trigger_type, created_at + FROM invocations + WHERE function_name = $1 AND namespace = $2 + ORDER BY created_at DESC + LIMIT $3 + `, functionName, namespace, limit) + if err != nil { + return nil, fmt.Errorf("list invocations: %w", err) + } + defer rows.Close() + + var result []*Invocation + for rows.Next() { + inv := &Invocation{} + if err := rows.Scan( + &inv.ID, &inv.FunctionName, &inv.Namespace, + &inv.Status, &inv.DurationMs, &inv.HTTPStatus, + &inv.Logs, &inv.TriggerType, &inv.CreatedAt, + ); err != nil { + return nil, fmt.Errorf("scan invocation: %w", err) + } + result = append(result, inv) + } + return result, rows.Err() +} + +// RunMigrations применяет SQL файлы из директории migrations. +// Простая реализация без внешних зависимостей — выполняем один файл. +// Используем IF NOT EXISTS в SQL поэтому безопасно запускать повторно. +func (s *Store) RunMigrations(ctx context.Context, sql string) error { + _, err := s.db.ExecContext(ctx, sql) + if err != nil { + return fmt.Errorf("run migration: %w", err) + } + return nil +} diff --git a/internal/storage/s3/.gitkeep b/internal/storage/s3/.gitkeep deleted file mode 100644 index e69de29..0000000 diff --git a/internal/storage/s3/client.go b/internal/storage/s3/client.go new file mode 100644 index 0000000..07d432c --- /dev/null +++ b/internal/storage/s3/client.go @@ -0,0 +1,82 @@ +// Изменено: 2026-03-07 +// S3 storage — загрузка и скачивание zip архивов с кодом функций. +// Используем Ceph S3 compatible API (minio-go клиент умеет работать с любым S3). +// Код загружается пользователем через REST API, хранится в S3, +// builder скачивает его для сборки Docker образа. + +package s3 + +import ( + "context" + "fmt" + "io" + + "github.com/minio/minio-go/v7" + "github.com/minio/minio-go/v7/pkg/credentials" +) + +// Client — клиент для работы с S3. +type Client struct { + mc *minio.Client + bucket string +} + +// New создаёт клиент S3 и проверяет/создаёт бакет. +func New(endpoint, accessKey, secretKey, bucket string, useSSL bool) (*Client, error) { + mc, err := minio.New(endpoint, &minio.Options{ + Creds: credentials.NewStaticV4(accessKey, secretKey, ""), + Secure: useSSL, + }) + if err != nil { + return nil, fmt.Errorf("create minio client: %w", err) + } + return &Client{mc: mc, bucket: bucket}, nil +} + +// EnsureBucket создаёт бакет если он не существует. +// Вызывается при старте сервиса. +func (c *Client) EnsureBucket(ctx context.Context) error { + exists, err := c.mc.BucketExists(ctx, c.bucket) + if err != nil { + return fmt.Errorf("check bucket: %w", err) + } + if !exists { + if err := c.mc.MakeBucket(ctx, c.bucket, minio.MakeBucketOptions{}); err != nil { + return fmt.Errorf("create bucket: %w", err) + } + } + return nil +} + +// Upload загружает zip архив с кодом функции в S3. +// Ключ: functions/{namespace}/{name}/{version}.zip +// Возвращает ключ объекта для сохранения в CRD. +func (c *Client) Upload(ctx context.Context, namespace, funcName, version string, r io.Reader, size int64) (string, error) { + key := fmt.Sprintf("functions/%s/%s/%s.zip", namespace, funcName, version) + _, err := c.mc.PutObject(ctx, c.bucket, key, r, size, minio.PutObjectOptions{ + ContentType: "application/zip", + }) + if err != nil { + return "", fmt.Errorf("upload function code: %w", err) + } + return key, nil +} + +// Download скачивает zip архив с кодом функции из S3. +// Используется builder'ом для сборки образа. +func (c *Client) Download(ctx context.Context, key string) (io.ReadCloser, error) { + obj, err := c.mc.GetObject(ctx, c.bucket, key, minio.GetObjectOptions{}) + if err != nil { + return nil, fmt.Errorf("download function code: %w", err) + } + return obj, nil +} + +// Delete удаляет архив с кодом функции из S3. +// Вызывается при удалении Function CRD. +func (c *Client) Delete(ctx context.Context, key string) error { + if err := c.mc.RemoveObject(ctx, c.bucket, key, minio.RemoveObjectOptions{}); err != nil { + return fmt.Errorf("delete function code: %w", err) + } + return nil +} diff --git a/main.go b/main.go index 774824f..0a6549c 100644 --- a/main.go +++ b/main.go @@ -1,23 +1,16 @@ -/* -Copyright 2026. - -Licensed under the Apache License, Version 2.0 (the "License"); -you may not use this file except in compliance with the License. -You may obtain a copy of the License at - - http://www.apache.org/licenses/LICENSE-2.0 - -Unless required by applicable law or agreed to in writing, software -distributed under the License is distributed on an "AS IS" BASIS, -WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -See the License for the specific language governing permissions and -limitations under the License. -*/ +// Изменено: 2026-03-07 +// main.go — точка входа. Запускает operator manager и REST API сервер параллельно. +// Operator manager управляет Function/Trigger CRD через reconcile loop. +// REST API (gorilla/mux) принимает запросы от Terraform provider. package main import ( + "context" "flag" + "fmt" + "log/slog" + "net/http" "os" // Import all Kubernetes client auth plugins (e.g. Azure, GCP, OIDC, etc.) @@ -33,12 +26,17 @@ import ( slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1" "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/controllers" + slessapi "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/api" + "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/api/handler" + "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/builder" + "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/config" + "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/storage/postgres" + "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/storage/s3" //+kubebuilder:scaffold:imports ) var ( - scheme = runtime.NewScheme() - setupLog = ctrl.Log.WithName("setup") + scheme = runtime.NewScheme() ) func init() { @@ -65,6 +63,46 @@ func main() { ctrl.SetLogger(zap.New(zap.UseFlagOptions(&opts))) + // slog для REST API (JSON-формат для prod, text для dev) + log := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo})) + + // Загружаем конфиг из env vars + cfg, err := config.Load() + if err != nil { + log.Error("load config", "err", err) + os.Exit(1) + } + + // Подключаемся к PostgreSQL + pg, err := postgres.New(cfg.PostgresDSN) + if err != nil { + log.Error("connect postgres", "err", err) + os.Exit(1) + } + defer pg.Close() + + // Применяем миграции БД — читаем SQL и передаём строкой + migrationSQL, err := os.ReadFile("migrations/001_initial.sql") + if err != nil { + log.Error("read migration file", "err", err) + os.Exit(1) + } + if err := pg.RunMigrations(context.Background(), string(migrationSQL)); err != nil { + log.Error("run migrations", "err", err) + os.Exit(1) + } + + // S3 клиент + s3Client, err := s3.New(cfg.S3Endpoint, cfg.S3AccessKey, cfg.S3SecretKey, cfg.S3Bucket, cfg.S3UseSSL) + if err != nil { + log.Error("connect s3", "err", err) + os.Exit(1) + } + if err := s3Client.EnsureBucket(context.Background()); err != nil { + log.Error("ensure s3 bucket", "err", err) + os.Exit(1) + } + mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), ctrl.Options{ Scheme: scheme, MetricsBindAddress: metricsAddr, @@ -72,51 +110,71 @@ func main() { HealthProbeBindAddress: probeAddr, LeaderElection: enableLeaderElection, LeaderElectionID: "4b6ba465.kube5s.ru", - // LeaderElectionReleaseOnCancel defines if the leader should step down voluntarily - // when the Manager ends. This requires the binary to immediately end when the - // Manager is stopped, otherwise, this setting is unsafe. Setting this significantly - // speeds up voluntary leader transitions as the new leader don't have to wait - // LeaseDuration time first. - // - // In the default scaffold provided, the program ends immediately after - // the manager stops, so would be fine to enable this option. However, - // if you are doing or is intended to do any operation such as perform cleanups - // after the manager stops then its usage might be unsafe. - // LeaderElectionReleaseOnCancel: true, }) if err != nil { - setupLog.Error(err, "unable to start manager") + log.Error("unable to start manager", "err", err) os.Exit(1) } + // Builder использует k8s client из manager'а + bldr := builder.New(mgr.GetClient(), builder.Config{ + BuilderImage: cfg.BuilderImage, + RegistryHost: cfg.RegistryHost, + S3Endpoint: cfg.S3Endpoint, + S3AccessKey: cfg.S3AccessKey, + S3SecretKey: cfg.S3SecretKey, + S3Bucket: cfg.S3Bucket, + Namespace: "sless", + }) + if err = (&controllers.FunctionReconciler{ - Client: mgr.GetClient(), - Scheme: mgr.GetScheme(), + Client: mgr.GetClient(), + Scheme: mgr.GetScheme(), + Builder: bldr, }).SetupWithManager(mgr); err != nil { - setupLog.Error(err, "unable to create controller", "controller", "Function") + log.Error("unable to create controller", "controller", "Function", "err", err) os.Exit(1) } if err = (&controllers.TriggerReconciler{ - Client: mgr.GetClient(), - Scheme: mgr.GetScheme(), + Client: mgr.GetClient(), + Scheme: mgr.GetScheme(), + IngressHost: cfg.IngressHost, }).SetupWithManager(mgr); err != nil { - setupLog.Error(err, "unable to create controller", "controller", "Trigger") + log.Error("unable to create controller", "controller", "Trigger", "err", err) os.Exit(1) } //+kubebuilder:scaffold:builder if err := mgr.AddHealthzCheck("healthz", healthz.Ping); err != nil { - setupLog.Error(err, "unable to set up health check") + log.Error("unable to set up health check", "err", err) os.Exit(1) } if err := mgr.AddReadyzCheck("readyz", healthz.Ping); err != nil { - setupLog.Error(err, "unable to set up ready check") + log.Error("unable to set up ready check", "err", err) os.Exit(1) } - setupLog.Info("starting manager") + // REST API сервер — запускается параллельно с operator manager + apiHandler := slessapi.NewRouter(&handler.Handler{ + K8s: mgr.GetClient(), + Scheme: mgr.GetScheme(), + S3: s3Client, + PG: pg, + Log: log, + }, cfg.APIToken, log) + + go func() { + addr := fmt.Sprintf(":%d", cfg.APIPort) + log.Info("starting REST API", "addr", addr) + if err := http.ListenAndServe(addr, apiHandler); err != nil { + log.Error("REST API server failed", "err", err) + os.Exit(1) + } + }() + + log.Info("starting operator manager") if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil { - setupLog.Error(err, "problem running manager") + log.Error("problem running manager", "err", err) os.Exit(1) } } diff --git a/migrations/001_initial.sql b/migrations/001_initial.sql new file mode 100644 index 0000000..84525a2 --- /dev/null +++ b/migrations/001_initial.sql @@ -0,0 +1,26 @@ +-- Миграция: 001 — начальная схема +-- Создано: 2026-03-07 +-- Таблица invocations хранит логи вызовов функций. +-- Состояние самих функций хранится в k8s CRD (etcd), не здесь. + +CREATE TABLE IF NOT EXISTS invocations ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + -- function_name и namespace — ссылка на CRD Function в k8s + function_name VARCHAR(255) NOT NULL, + namespace VARCHAR(255) NOT NULL, + -- статус выполнения: success, error, timeout + status VARCHAR(50) NOT NULL, + -- время выполнения в миллисекундах + duration_ms INTEGER NOT NULL DEFAULT 0, + -- HTTP статус ответа функции (для HTTP триггеров) + http_status INTEGER, + -- вывод функции (stdout/stderr), ограничен 64KB + logs TEXT, + -- тип триггера: http, cron + trigger_type VARCHAR(50) NOT NULL DEFAULT 'http', + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); + +-- Индексы для быстрой выборки логов по функции +CREATE INDEX IF NOT EXISTS idx_invocations_function ON invocations (function_name, namespace); +CREATE INDEX IF NOT EXISTS idx_invocations_created_at ON invocations (created_at DESC);