diff --git a/controllers/function_controller.go b/controllers/function_controller.go index 2da300b..5b4f7a1 100644 --- a/controllers/function_controller.go +++ b/controllers/function_controller.go @@ -106,11 +106,36 @@ func (r *FunctionReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c const finalizerName = "sless.kube5s.ru/finalizer" -// startBuild запускает kaniko Job и помечает функцию как Building. +// startBuild проверяет наличие образа в registry и либо пропускает сборку, +// либо запускает kaniko Job. Идемпотентность: если код не менялся (тег = hash s3Key), +// образ уже в registry → deploy без пересборки. // Критически важно: СНАЧАЛА сохраняем last-built-s3key аннотацию, ПОТОМ status. -// Это предотвращает повторный запуск сборки при параллельных reconcile — -// следующий reconcile увидит last-built-s3key == spec.S3Key и не войдёт в startBuild. func (r *FunctionReconciler) startBuild(ctx context.Context, fn *slessv1alpha1.Function) (ctrl.Result, error) { + imageRef := r.Builder.ImageRef(fn.Namespace, fn.Name, fn.Spec.S3Key) + + // Проверяем: образ с этим тегом уже существует в registry? + // Если да — пропускаем kaniko, сразу переходим в Ready. + if r.Builder.ImageExists(ctx, imageRef) { + logger := log.FromContext(ctx) + logger.Info("image already exists in registry, skipping build", "imageRef", imageRef) + + if fn.Annotations == nil { + fn.Annotations = map[string]string{} + } + fn.Annotations["sless.kube5s.ru/last-built-s3key"] = fn.Spec.S3Key + if err := r.Update(ctx, fn); err != nil { + return ctrl.Result{}, fmt.Errorf("update annotations (cache hit): %w", err) + } + + fn.Status.Phase = slessv1alpha1.FunctionPhaseReady + fn.Status.ImageRef = imageRef + fn.Status.Message = "Image restored from registry cache" + if err := r.Status().Update(ctx, fn); err != nil { + return ctrl.Result{}, fmt.Errorf("update status (cache hit): %w", err) + } + return ctrl.Result{}, nil + } + 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)) diff --git a/controllers/functionjob_controller.go b/controllers/functionjob_controller.go index 761257b..ee2d11e 100644 --- a/controllers/functionjob_controller.go +++ b/controllers/functionjob_controller.go @@ -109,7 +109,8 @@ func (r *FunctionJobReconciler) Reconcile(ctx context.Context, req ctrl.Request) return r.createRunJob(ctx, fj, deployNS, jobName) } -// startJobBuild запускает kaniko сборку образа и переводит FunctionJob в фазу Building. +// startJobBuild проверяет наличие образа в registry и либо пропускает сборку, +// либо запускает kaniko Job. Идемпотентность по hash s3Key аналогична service/function. func (r *FunctionJobReconciler) startJobBuild(ctx context.Context, fj *slessv1alpha1.FunctionJob) (ctrl.Result, error) { logger := log.FromContext(ctx) @@ -120,6 +121,28 @@ func (r *FunctionJobReconciler) startJobBuild(ctx context.Context, fj *slessv1al return ctrl.Result{RequeueAfter: 5 * time.Second}, nil } + imageRef := r.Builder.ImageRef(r.OperatorNamespace, fj.Name, fj.Spec.S3Key) + + // Проверяем: образ с этим тегом уже существует в registry? + if r.Builder.ImageExists(ctx, imageRef) { + logger.Info("image already exists in registry, skipping build", "imageRef", imageRef) + + if fj.Annotations == nil { + fj.Annotations = map[string]string{} + } + if err := r.Update(ctx, fj); err != nil { + return ctrl.Result{}, fmt.Errorf("update annotations (cache hit): %w", err) + } + + fj.Status.Phase = slessv1alpha1.FunctionJobPhaseBuilding // перейдёт в run сразу + fj.Status.ImageRef = imageRef + fj.Status.Message = "Image restored from registry cache" + if err := r.Status().Update(ctx, fj); err != nil { + return ctrl.Result{}, fmt.Errorf("update status (cache hit): %w", err) + } + return ctrl.Result{Requeue: true}, nil + } + buildJobName, err := r.Builder.Build(ctx, r.OperatorNamespace, fj.Name, fj.Spec.S3Key) if err != nil { fj.Status.Phase = slessv1alpha1.FunctionJobPhaseFailed diff --git a/controllers/service_controller.go b/controllers/service_controller.go index a45d93b..d9a77b3 100644 --- a/controllers/service_controller.go +++ b/controllers/service_controller.go @@ -101,8 +101,35 @@ func (r *ServiceReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ct return ctrl.Result{}, nil } -// startServiceBuild запускает kaniko Job и помечает сервис как Building. +// startServiceBuild проверяет наличие образа в registry и либо пропускает сборку, +// либо запускает kaniko Job. Идемпотентность: если код не менялся (тег = hash s3Key), +// образ уже в registry → deploy без пересборки. func (r *ServiceReconciler) startServiceBuild(ctx context.Context, svc *slessv1alpha1.Service) (ctrl.Result, error) { + imageRef := r.Builder.ImageRef(svc.Namespace, svc.Name, svc.Spec.S3Key) + + // Проверяем: образ с этим тегом уже существует в registry? + // Если да — пропускаем kaniko, сразу переходим в Ready с известным imageRef. + if r.Builder.ImageExists(ctx, imageRef) { + logger := log.FromContext(ctx) + logger.Info("image already exists in registry, skipping build", "imageRef", imageRef) + + if svc.Annotations == nil { + svc.Annotations = map[string]string{} + } + svc.Annotations["sless.kube5s.ru/last-built-s3key"] = svc.Spec.S3Key + if err := r.Update(ctx, svc); err != nil { + return ctrl.Result{}, fmt.Errorf("update service annotations (cache hit): %w", err) + } + + svc.Status.Phase = slessv1alpha1.ServicePhaseReady + svc.Status.ImageRef = imageRef + svc.Status.Message = "Image restored from registry cache" + if err := r.Status().Update(ctx, svc); err != nil { + return ctrl.Result{}, fmt.Errorf("update service status (cache hit): %w", err) + } + return ctrl.Result{}, nil + } + jobName, err := r.Builder.Build(ctx, svc.Namespace, svc.Name, svc.Spec.S3Key) if err != nil { return r.setServiceFailed(ctx, svc, fmt.Sprintf("failed to start build: %v", err)) diff --git a/doc/errors/log.md b/doc/errors/log.md index 62f18cd..3cf97e7 100644 --- a/doc/errors/log.md +++ b/doc/errors/log.md @@ -4,6 +4,68 @@ --- +## 2026-03-22 — БАГ: CreateService/CreateFunction возвращает 409 при `terraform apply -replace` (ИСПРАВЛЕН) + +### Симптом + +``` +terraform apply -replace=sless_service.pg_info + sless_service.pg_info: Destroying... [name=pg-info] + sless_service.pg_info: Destruction complete + sless_service.pg_info: Creating... + Error: create service: status 409: {"error":"service already exists"} +``` + +Все 22+ сервиса из `-replace` падают с 409 при пересоздании. + +### Точная причина + +**Цепочка событий** (воспроизводится только при наличии finalizer): + +1. `terraform` вызывает `DELETE /services/pg-info` +2. API делает `h.K8s.Delete(svc)` → k8s **НЕ удаляет объект** немедленно. + Вместо этого — выставляет `DeletionTimestamp` на объекте и ждёт. +3. `service_controller.go` (асинхронно!) обрабатывает удаление: + - сносит Deployment, k8s Service, Ingress + - затем вызывает `svc.Finalizers = removeString(svc.Finalizers, serviceFinalizerName)` + - только после этого k8s реально удаляет CRD объект из etcd + - **латентность: 1–5 секунд** +4. `terraform` **немедленно** вызывает `POST /services/pg-info` +5. API делает `h.K8s.Create(svc)` → etcd возвращает `IsAlreadyExists` +6. Старый код проверял: `phase == Failed`? → нет, было `Ready` → `shouldRecreate=false` → **409** + +**Ключевое: объект существует с `DeletionTimestamp != zero`, то есть он уже "мёртвый", но finalizer ещё не снят. Старый код этого не проверял.** + +### Файл с багом + +`internal/api/handler/services.go` → `CreateService()` — строка `shouldRecreate` +`internal/api/handler/functions.go` → `CreateFunction()` — аналогичная логика + +### Исправление + +В обоих файлах добавлена проверка `!existing.DeletionTimestamp.IsZero()` **до** проверки `shouldRecreate`: + +```go +if getErr == nil && !existing.DeletionTimestamp.IsZero() { + // Объект удаляется: polling каждую секунду до 30 сек + for i := 0; i < 30; i++ { + time.Sleep(1 * time.Second) + if errors.IsNotFound(h.K8s.Get(...)) { + // исчez — создаём + } + } + // таймаут — возвращаем 409 с "try again later" +} +``` + +### Почему именно polling, а не Watch + +`http.ResponseWriter` не поддерживает long-poll без дополнительного механизма. +Watch на объект внутри HTTP handler — антипаттерн (использует goroutine leak при отмене). +30 × 1s — достаточно для любого разумного кластера; terraform имеет свой retry. + +--- + ## 2026-03-22 — ПОВЕДЕНИЕ: оператор выставляет Ready до готовности pod (known limitation) ### Симптом diff --git a/doc/progress.md b/doc/progress.md index 94c5fc7..c217278 100644 --- a/doc/progress.md +++ b/doc/progress.md @@ -1,6 +1,30 @@ # Прогресс разработки -Последнее обновление: 2026-03-22 +Последнее обновление: 2026-03-23 + +--- + +## 2026-03-23 — Сессия 10: ImageExists + деплой v0.1.58 + +### Что сделано + +- **ImageExists** — новый метод `Builder.ImageExists(ctx, imageRef) bool` в `internal/builder/builder.go` + - Docker Registry v2 API: anonymous bearer token → HEAD `/v2/{repo}/manifests/{tag}` + - timeout 5s, fallback false при любой ошибке (безопасный fallback → kaniko) + - Приватные registry (Harbor и др.) → 401 → false → строим заново +- **service_controller.go** `startServiceBuild()`: кеш-проверка перед Build(); cache hit → `Status.Phase=Ready` + `Status.ImageRef` без kaniko +- **function_controller.go** `startBuild()`: аналогично +- **functionjob_controller.go** `startJobBuild()`: аналогично; cache hit → Requeue для немедленного перехода к run +- **v0.1.58** собран и задеплоен: `naeel/sless-operator:v0.1.58` Running, RESTARTS:0 + +### Статус +✅ Код написан (no errors) +✅ Docker build + push наDockerHub (sha256:7b0d48a01a29f2a5...) +✅ kubectl rollout: `deployment "sless-operator" successfully rolled out` +✅ Pod `sless-operator-85d4685c4b-6nxqw` Running v0.1.58 + +### Следующий шаг +Протестировать cache hit: вернуть `.tf.disabled` → `.tf`, `terraform apply` — образы уже на DockerHub, должны восстановиться без kaniko за секунды. --- diff --git a/examples/POSTGRES/chaos_marathon.tf b/examples/POSTGRES/chaos_marathon.tf index 6fdb035..88ee5a7 100644 --- a/examples/POSTGRES/chaos_marathon.tf +++ b/examples/POSTGRES/chaos_marathon.tf @@ -1,5 +1,5 @@ // 2026-03-21 — chaos_marathon.tf: 15 новых сервисов для часового хаос-марафона. -// Три рантайма: python3.11 (9), nodejs20 (2), go1.23 (2). +// Два рантайма: python3.11 (9), nodejs20 (2). // Все зависят от sless_job.postgres_table_init_job. # ── Python: работа с таблицей ───────────────────────────────────────────────── @@ -233,50 +233,6 @@ resource "sless_service" "js_idempotent" { depends_on = [sless_job.postgres_table_init_job] } -# ── Go 1.23 ─────────────────────────────────────────────────────────────────── - -# Параллельные concurrent INSERTs из N горутин внутри одного пода. -resource "sless_service" "go_pg_race" { - name = "go-pg-race" - runtime = "go1.23" - entrypoint = "handler.Handle" - memory_mb = 256 - timeout_sec = 30 - - env_vars = { - PGHOST = local.pg_host - PGPORT = "5432" - PGDATABASE = local.pg_database - PGUSER = local.pg_username - PGPASSWORD = local.pg_password - PGSSLMODE = "require" - } - - source_dir = "${path.module}/code/go-pg-race" - depends_on = [sless_job.postgres_table_init_job] -} - -# Atomic-счётчик в памяти + PG INSERT на каждый вызов. -resource "sless_service" "go_counter_atomic" { - name = "go-counter-atomic" - runtime = "go1.23" - entrypoint = "handler.Handle" - memory_mb = 128 - timeout_sec = 15 - - env_vars = { - PGHOST = local.pg_host - PGPORT = "5432" - PGDATABASE = local.pg_database - PGUSER = local.pg_username - PGPASSWORD = local.pg_password - PGSSLMODE = "require" - } - - source_dir = "${path.module}/code/go-counter-atomic" - depends_on = [sless_job.postgres_table_init_job] -} - # ── Python: retry ───────────────────────────────────────────────────────────── # Запись с retry при transient PG error — тест устойчивости к сбоям. diff --git a/examples/POSTGRES/code/go-counter-atomic/handler.go b/examples/POSTGRES/code/go-counter-atomic/handler.go deleted file mode 100644 index e101ab3..0000000 --- a/examples/POSTGRES/code/go-counter-atomic/handler.go +++ /dev/null @@ -1,55 +0,0 @@ -// 2026-03-21 — go-counter-atomic: считает вызовы через atomic в памяти + пишет в PG. -// Тестирует: in-memory state между вызовами (Go pod остаётся живым), + PG INSERT на каждый вызов. -package handler - -import ( - "context" - "fmt" - "os" - "sync/atomic" - "time" - - "github.com/jackc/pgx/v5/pgxpool" -) - -// invocations считается между вызовами (пока pod жив). -var invocations int64 - -// Handle записывает факт вызова в PG и возвращает накопленный счётчик. -func Handle(event map[string]interface{}) interface{} { - n := atomic.AddInt64(&invocations, 1) - - dsn := fmt.Sprintf( - "host=%s port=%s dbname=%s user=%s password=%s sslmode=%s", - os.Getenv("PGHOST"), envOrDefault("PGPORT", "5432"), - os.Getenv("PGDATABASE"), os.Getenv("PGUSER"), - os.Getenv("PGPASSWORD"), envOrDefault("PGSSLMODE", "require"), - ) - pool, err := pgxpool.New(context.Background(), dsn) - if err != nil { - return map[string]interface{}{"invocation_n": n, "error": err.Error()} - } - defer pool.Close() - - title := fmt.Sprintf("go-counter-invoke-%d-%d", n, time.Now().UnixMilli()) - var id int64 - err = pool.QueryRow(context.Background(), - "INSERT INTO terraform_demo_table (title) VALUES ($1) RETURNING id", title, - ).Scan(&id) - if err != nil { - return map[string]interface{}{"invocation_n": n, "error": err.Error()} - } - - return map[string]interface{}{ - "invocation_n": n, - "inserted_id": id, - "title": title, - } -} - -func envOrDefault(key, def string) string { - if v := os.Getenv(key); v != "" { - return v - } - return def -} diff --git a/examples/POSTGRES/code/go-pg-race/handler.go b/examples/POSTGRES/code/go-pg-race/handler.go deleted file mode 100644 index 4dd1b39..0000000 --- a/examples/POSTGRES/code/go-pg-race/handler.go +++ /dev/null @@ -1,94 +0,0 @@ -// 2026-03-21 — go-pg-race: параллельные INSERT из нескольких горутин внутри одной функции. -// Тестирует: race condition устойчивость Go + PG при concurrent writes из одного пода. -// Использует pgx/v5 (pre-cached в base image). -package handler - -import ( - "context" - "fmt" - "os" - "sync" - "sync/atomic" - "time" - - "github.com/jackc/pgx/v5/pgxpool" -) - -// Handle запускает workers горутин, каждая делает n_per_worker INSERTs. -func Handle(event map[string]interface{}) interface{} { - workers := intParam(event, "workers", 5) - if workers > 20 { - workers = 20 - } - nPerWorker := intParam(event, "n_per_worker", 10) - if nPerWorker > 50 { - nPerWorker = 50 - } - - dsn := fmt.Sprintf( - "host=%s port=%s dbname=%s user=%s password=%s sslmode=%s", - os.Getenv("PGHOST"), getenv("PGPORT", "5432"), - os.Getenv("PGDATABASE"), os.Getenv("PGUSER"), - os.Getenv("PGPASSWORD"), getenv("PGSSLMODE", "require"), - ) - pool, err := pgxpool.New(context.Background(), dsn) - if err != nil { - return map[string]interface{}{"error": err.Error()} - } - defer pool.Close() - - var ( - wg sync.WaitGroup - ok int64 - errCount int64 - ) - t0 := time.Now() - for w := 0; w < workers; w++ { - wg.Add(1) - go func(wid int) { - defer wg.Done() - for i := 0; i < nPerWorker; i++ { - title := fmt.Sprintf("go-race-w%d-%d-%d", wid, time.Now().UnixMilli(), i) - _, err := pool.Exec(context.Background(), - "INSERT INTO terraform_demo_table (title) VALUES ($1)", title) - if err != nil { - atomic.AddInt64(&errCount, 1) - } else { - atomic.AddInt64(&ok, 1) - } - } - }(w) - } - wg.Wait() - elapsed := time.Since(t0).Seconds() - - return map[string]interface{}{ - "workers": workers, - "n_per_worker": nPerWorker, - "inserted": ok, - "errors": errCount, - "elapsed_sec": elapsed, - "ops_per_sec": float64(ok) / elapsed, - } -} - -func intParam(event map[string]interface{}, key string, def int) int { - v, ok := event[key] - if !ok { - return def - } - switch val := v.(type) { - case float64: - return int(val) - case int: - return val - } - return def -} - -func getenv(key, def string) string { - if v := os.Getenv(key); v != "" { - return v - } - return def -} diff --git a/examples/POSTGRES/code/stress-go-fast/handler.go b/examples/POSTGRES/code/stress-go-fast/handler.go deleted file mode 100644 index 9f40709..0000000 --- a/examples/POSTGRES/code/stress-go-fast/handler.go +++ /dev/null @@ -1,42 +0,0 @@ -// 2026-03-19 -// handler.go — быстрая Go функция: факториал + числа Фибоначчи. -// Проверяет Go runtime под лёгкой нагрузкой и корректность JSON-ответа. -// Entrypoint: handler.Handle -package handler - -import "fmt" - -func factorial(n int) uint64 { - if n <= 1 { - return 1 - } - return uint64(n) * factorial(n-1) -} - -func fib(n int) int { - if n <= 1 { - return n - } - a, b := 0, 1 - for i := 2; i <= n; i++ { - a, b = b, a+b - } - return b -} - -func Handle(event map[string]interface{}) interface{} { - n := 10 - if v, ok := event["n"].(float64); ok { - n = int(v) - if n > 20 { - n = 20 - } - } - return map[string]interface{}{ - "runtime": "go1.23", - "version": "v1", - "n": n, - "factorial": fmt.Sprintf("%d", factorial(n)), - "fib": fib(n), - } -} diff --git a/examples/POSTGRES/code/stress-go-nil/handler.go b/examples/POSTGRES/code/stress-go-nil/handler.go deleted file mode 100644 index 0ee25f6..0000000 --- a/examples/POSTGRES/code/stress-go-nil/handler.go +++ /dev/null @@ -1,21 +0,0 @@ -// 2026-03-19 -// handler.go — намеренный nil pointer dereference в Go. -// Проверяет что Go runtime recover() перехватывает панику и платформа возвращает 500. -// Entrypoint: handler.Handle -package handler - -func Handle(event map[string]interface{}) interface{} { - crash := true - if v, ok := event["crash"].(bool); ok { - crash = v - } - if crash { - var p *string - _ = *p // panic: намеренный nil pointer для stress-теста - } - return map[string]interface{}{ - "runtime": "go1.23", - "version": "v1", - "crashed": false, - } -} diff --git a/examples/POSTGRES/code/stress-go-pgstorm/handler.go b/examples/POSTGRES/code/stress-go-pgstorm/handler.go deleted file mode 100644 index edf4668..0000000 --- a/examples/POSTGRES/code/stress-go-pgstorm/handler.go +++ /dev/null @@ -1,148 +0,0 @@ -// 2026-03-19 -// handler.go — Go стресс-тест PostgreSQL через pgxpool. -// Запускает N горутин (default 100), каждая в цикле duration_sec (default 600) -// долбит PG попеременно: INSERT / SELECT COUNT / SELECT MAX с случайными задержками. -// Цель: проверить Go runtime под конкурентной нагрузкой и устойчивость PG connection pool. -// Entrypoint: handler.Handle -package handler - -import ( - "context" - "fmt" - "math/rand" - "os" - "sync" - "sync/atomic" - "time" - - "github.com/jackc/pgx/v5/pgxpool" -) - -// pgDSN собирает DSN из env vars (PGHOST, PGPORT, PGDATABASE, PGUSER, PGPASSWORD, PGSSLMODE). -func pgDSN() string { - host := os.Getenv("PGHOST") - port := os.Getenv("PGPORT") - if port == "" { - port = "5432" - } - db := os.Getenv("PGDATABASE") - user := os.Getenv("PGUSER") - pass := os.Getenv("PGPASSWORD") - sslmode := os.Getenv("PGSSLMODE") - if sslmode == "" { - sslmode = "require" - } - return fmt.Sprintf("host=%s port=%s dbname=%s user=%s password=%s sslmode=%s", - host, port, db, user, pass, sslmode) -} - -// worker — одна горутина: чередует INSERT/COUNT/MAX с случайной задержкой до maxDelayMs. -// При ошибке инкрементирует errOps и продолжает (не паникует). -func worker(ctx context.Context, pool *pgxpool.Pool, workerID int, maxDelayMs int, okOps, errOps *int64) { - rng := rand.New(rand.NewSource(time.Now().UnixNano() + int64(workerID))) - op := 0 - for { - select { - case <-ctx.Done(): - return - default: - } - - // Случайная задержка перед следующей операцией: 0..maxDelayMs мс - delay := rng.Intn(maxDelayMs + 1) - time.Sleep(time.Duration(delay) * time.Millisecond) - - var err error - switch op % 3 { - case 0: // INSERT - title := fmt.Sprintf("pgstorm-w%d-%d", workerID, time.Now().UnixNano()) - _, err = pool.Exec(ctx, - "INSERT INTO terraform_demo_table (title) VALUES ($1)", title) - case 1: // SELECT COUNT - var count int64 - err = pool.QueryRow(ctx, - "SELECT COUNT(*) FROM terraform_demo_table").Scan(&count) - case 2: // SELECT MAX id - var maxID *int64 - err = pool.QueryRow(ctx, - "SELECT MAX(id) FROM terraform_demo_table").Scan(&maxID) - } - - if err != nil && ctx.Err() == nil { - atomic.AddInt64(errOps, 1) - } else if err == nil { - atomic.AddInt64(okOps, 1) - } - op++ - } -} - -func Handle(event map[string]interface{}) interface{} { - // Параметры из event (все опциональны — разумные defaults) - workers := 100 - if v, ok := event["workers"].(float64); ok && v > 0 && v <= 500 { - workers = int(v) - } - durationSec := 600 - if v, ok := event["duration_sec"].(float64); ok && v > 0 && v <= 3600 { - durationSec = int(v) - } - maxDelayMs := 300 - if v, ok := event["max_delay_ms"].(float64); ok && v >= 0 && v <= 5000 { - maxDelayMs = int(v) - } - - // Инициализация pgxpool — единый pool на всю функцию, MaxConns ограничен - // чтобы не перегрузить managed PG при большом числе горутин. - poolCfg, err := pgxpool.ParseConfig(pgDSN()) - if err != nil { - return map[string]interface{}{"error": fmt.Sprintf("parse dsn: %v", err)} - } - maxConns := 20 - if workers < 20 { - maxConns = workers - } - poolCfg.MaxConns = int32(maxConns) - - ctx, cancel := context.WithTimeout(context.Background(), time.Duration(durationSec)*time.Second) - defer cancel() - - pool, err := pgxpool.NewWithConfig(ctx, poolCfg) - if err != nil { - return map[string]interface{}{"error": fmt.Sprintf("connect pool: %v", err)} - } - defer pool.Close() - - var okOps, errOps int64 - startTime := time.Now() - - var wg sync.WaitGroup - for i := 0; i < workers; i++ { - wg.Add(1) - go func(id int) { - defer wg.Done() - worker(ctx, pool, id, maxDelayMs, &okOps, &errOps) - }(i) - } - wg.Wait() - - elapsed := time.Since(startTime).Seconds() - total := okOps + errOps - opsPerSec := 0.0 - if elapsed > 0 { - opsPerSec = float64(total) / elapsed - } - - return map[string]interface{}{ - "runtime": "go1.23", - "version": "v1", - "workers": workers, - "duration_sec": durationSec, - "max_delay_ms": maxDelayMs, - "elapsed_sec": fmt.Sprintf("%.1f", elapsed), - "total_ops": total, - "ok_ops": okOps, - "err_ops": errOps, - "ops_per_sec": fmt.Sprintf("%.1f", opsPerSec), - } -} diff --git a/examples/POSTGRES/stress.tf b/examples/POSTGRES/stress.tf index b31d449..fc3281e 100644 --- a/examples/POSTGRES/stress.tf +++ b/examples/POSTGRES/stress.tf @@ -1,55 +1,7 @@ // 2026-03-21 — stress.tf: все стресс-сервисы для комплексного тестирования. -// Три рантайма: go1.23 (3), nodejs20 (2), python3.11 (5). +// Два рантайма: nodejs20 (2), python3.11 (5). // Все depends_on = [sless_job.postgres_table_init_job] — таблица должна существовать. -# ── Go 1.23 ───────────────────────────────────────────────────────────────── - -# Быстрая математика: факториал + числа Фибоначчи. Без PG. Проверяет Go runtime. -resource "sless_service" "stress_go_fast" { - name = "stress-go-fast" - runtime = "go1.23" - entrypoint = "handler.Handle" - memory_mb = 128 - timeout_sec = 15 - source_dir = "${path.module}/code/stress-go-fast" - - depends_on = [sless_job.postgres_table_init_job] -} - -# Намеренный nil pointer dereference. Без PG. Проверяет recover() в Go runtime. -resource "sless_service" "stress_go_nil" { - name = "stress-go-nil" - runtime = "go1.23" - entrypoint = "handler.Handle" - memory_mb = 128 - timeout_sec = 10 - source_dir = "${path.module}/code/stress-go-nil" - - depends_on = [sless_job.postgres_table_init_job] -} - -# Конкурентный PG-шторм через pgxpool: N горутин INSERT/COUNT/MAX. -# timeout_sec=660 — покрывает max duration_sec=600 с запасом. -# pgx/v5 уже в базовом образе go1.23: user go.mod не нужен. -resource "sless_service" "stress_go_pgstorm" { - name = "stress-go-pgstorm" - runtime = "go1.23" - entrypoint = "handler.Handle" - memory_mb = 256 - timeout_sec = 660 - source_dir = "${path.module}/code/stress-go-pgstorm" - - env_vars = { - PGHOST = local.pg_host - PGPORT = "5432" - PGDATABASE = local.pg_database - PGUSER = local.pg_username - PGPASSWORD = local.pg_password - PGSSLMODE = "require" - } - - depends_on = [sless_job.postgres_table_init_job] -} # ── Node.js 20 ──────────────────────────────────────────────────────────────── diff --git a/internal/api/handler/functions.go b/internal/api/handler/functions.go index 8d1d25b..69d6200 100644 --- a/internal/api/handler/functions.go +++ b/internal/api/handler/functions.go @@ -1,5 +1,6 @@ // Изменено: 2026-03-18 (добавлены created_at, last_built_at в functionResponse и fnToResponse) // Изменено: 2026-03-21 (fix: DeleteFunction возвращает 404 вместо 204 при отсутствующем объекте) +// Изменено: 2026-03-22 (fix: CreateFunction 409 при пересоздании функции — не учитывался DeletionTimestamp) // functions.go — CRUD handlers для Function CRD. // Принимает JSON, создаёт/обновляет/удаляет k8s ресурсы Function. // Namespace берётся из URL: /v1/namespaces/{namespace}/functions/{name} @@ -9,6 +10,7 @@ package handler import ( "encoding/json" "net/http" + "time" "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -128,22 +130,59 @@ func (h *Handler) CreateFunction(w http.ResponseWriter, r *http.Request) { S3Key: req.S3Key, }, } + // Пытаемся создать Function CRD в k8s (запись в etcd через controller-runtime). if err := h.K8s.Create(r.Context(), fn); err != nil { + + // IsAlreadyExists = etcd вернул 409 CONFLICT. + // Возникает в двух основных сценариях: + // + // [A] terraform apply -replace (= delete + create за один apply): + // 1. terraform DELETE /functions/{name} → API вызывает h.K8s.Delete(fn) + // 2. k8s ставит DeletionTimestamp и ждёт снятия finalizer sless.kube5s.ru/finalizer + // function_controller.go убивает kaniko Job и снимает finalizer асинхронно + // 3. terraform сразу POST /functions/{name} → IsAlreadyExists ← БАГ до этого фикса + // → НОВЫЙ КОД: обнаруживает DeletionTimestamp → ждёт исчезновения → создаёт + // + // [B] Split-brain кеша controller-runtime: + // Объект удалён из etcd, но informer-кеш ещё не обновился. + // Create падает с IsAlreadyExists из кеша, Get возвращает NotFound → пересоздаём. + // + // [C] Функция в фазе Failed: + // Предыдущий build провалился. Новый apply пытается создать снова. + // Удаляем Failed объект и пересоздаём. if errors.IsAlreadyExists(err) { - // IsAlreadyExists может прийти из кеша controller-runtime (split-brain): - // объект удалён из etcd, но кеш informer ещё южив. Делаем uncached Get: - // если реально NotFound — кеш устарел, пересоздаём. - // если существует и фаза Failed — тоже пересоздаём (build провалился, терраформ не добавил в state). - // если существует и фаза Ready/Building — возвращаем 409 (функция реально есть). + + // Получаем актуальное состояние объекта из etcd (не из кеша informer). existing := &slessv1alpha1.Function{} getErr := h.K8s.Get(r.Context(), client.ObjectKey{Name: req.Name, Namespace: ns}, existing) - shouldRecreate := errors.IsNotFound(getErr) || - (getErr == nil && existing.Status.Phase == slessv1alpha1.FunctionPhaseFailed) - if shouldRecreate { - if getErr == nil { - _ = h.K8s.Delete(r.Context(), existing) + + // --- Сценарий A: объект ожидает удаления --- + // DeletionTimestamp ≠ zero = k8s принял DELETE, finalizer ещё не снят. + // function_controller.go снимает finalizer после cleanup Job — обычно 1-5 сек. + if getErr == nil && !existing.DeletionTimestamp.IsZero() { + + deleted := false + // Polling каждую секунду, максимум 30 раз (= 30 секунд). + // 30 сек — запас на медленный кластер; в норме 1-3 итерации. + for i := 0; i < 30; i++ { + time.Sleep(1 * time.Second) + checkErr := h.K8s.Get(r.Context(), client.ObjectKey{Name: req.Name, Namespace: ns}, existing) + // NotFound = finalizer снят, объект исчез из etcd — можно создавать. + if errors.IsNotFound(checkErr) { + deleted = true + break + } + // Любая другая ошибка (timeout, сбой API) — продолжаем ждать. } - // Сбрасываем ResourceVersion — при split-brain etcd считает объект новым + + // Таймаут: объект не исчез за 30 секунд. + // Клиент (terraform) получит 409 и должен сделать retry позже. + if !deleted { + writeJSON(w, http.StatusConflict, errResp("function is being deleted, try again later")) + return + } + + // Объект исчез — сбрасываем ResourceVersion и создаём как новый. fn.ResourceVersion = "" if createErr := h.K8s.Create(r.Context(), fn); createErr != nil { writeJSON(w, http.StatusInternalServerError, errResp(createErr.Error())) @@ -152,6 +191,29 @@ func (h *Handler) CreateFunction(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusCreated, fnToResponse(fn)) return } + + // --- Сценарии B и C: split-brain или Failed --- + // NotFound при Get = кеш врёт (B); Failed фаза = сломанный build (C). + shouldRecreate := errors.IsNotFound(getErr) || + (getErr == nil && existing.Status.Phase == slessv1alpha1.FunctionPhaseFailed) + if shouldRecreate { + // Если объект реально есть (Failed) — сначала удаляем. + // Ошибку Delete игнорируем: Create покажет ошибку сам если что-то пошло не так. + if getErr == nil { + _ = h.K8s.Delete(r.Context(), existing) + } + // Сбрасываем ResourceVersion — при split-brain etcd считает объект новым. + fn.ResourceVersion = "" + if createErr := h.K8s.Create(r.Context(), fn); createErr != nil { + writeJSON(w, http.StatusInternalServerError, errResp(createErr.Error())) + return + } + writeJSON(w, http.StatusCreated, fnToResponse(fn)) + return + } + + // Объект живой (Ready/Building), DeletionTimestamp=zero. + // Легитимный конфликт — клиент пытается создать дубликат. writeJSON(w, http.StatusConflict, errResp("function already exists")) return } diff --git a/internal/api/handler/services.go b/internal/api/handler/services.go index c829b1e..a9df160 100644 --- a/internal/api/handler/services.go +++ b/internal/api/handler/services.go @@ -1,5 +1,6 @@ // Создано: 2026-03-20 (function-service-split) // Изменено: 2026-03-21 (fix: DeleteService возвращает 404 вместо 204 при отсутствующем объекте) +// Изменено: 2026-03-22 (fix: CreateService 409 при пересоздании сервиса через terraform -replace) // services.go — CRUD handlers для Service CRD (sless_service). // sless_service = long-running Deployment + URL. Каждый вызов проксируется к поду. // Namespace берётся из URL: /v1/namespaces/{namespace}/services/{name} @@ -134,18 +135,69 @@ func (h *Handler) CreateService(w http.ResponseWriter, r *http.Request) { S3Key: req.S3Key, }, } + // Пытаемся создать Service CRD в k8s. + // h.K8s.Create обращается к etcd через controller-runtime client. if err := h.K8s.Create(r.Context(), svc); err != nil { + + // IsAlreadyExists = etcd вернул 409 CONFLICT. + // Это нормально при конкурентных запросах или при "replace" через terraform. + // Сценарии: + // + // [A] terraform apply -replace: + // 1. terraform DELETE /services/{name} → API вызывает h.K8s.Delete(svc) + // 2. k8s НЕ удаляет объект немедленно: ставит DeletionTimestamp и ждёт + // пока service_controller.go снимет finalizer sless.kube5s.ru/service-finalizer + // (контроллер сначала сносит Deployment/Service/Ingress) + // 3. terraform сразу POST /services/{name} → h.K8s.Create → IsAlreadyExists + // (объект ещё есть в etcd, просто помечен на удаление) + // → СТАРЫЙ КОД: видел phase=Ready, shouldRecreate=false → возвращал 409 ← БАГ + // → НОВЫЙ КОД: видит DeletionTimestamp → ждёт исчезновения → создаёт + // + // [B] Кеш controller-runtime (split-brain): + // Объект удалён из etcd, но informer-кеш ещё не обновился. + // h.K8s.Create падает с IsAlreadyExists из кеша, хотя в etcd объекта нет. + // → Get возвращает NotFound → пересоздаём. + // + // [C] Сервис в фазе Failed: + // Предыдущий build провалился, terraform не добавил в state. + // Новый apply пытается создать снова → Failed объект удаляем и пересоздаём. if errors.IsAlreadyExists(err) { - // Повторяем логику function_handler: обрабатываем split-brain кеша. - // Если объект реально есть и не в фазе Failed — возвращаем 409. + + // Получаем актуальное состояние объекта напрямую из etcd (не из кеша). + // Нужно чтобы точно определить сценарий A/B/C. existing := &slessv1alpha1.Service{} getErr := h.K8s.Get(r.Context(), client.ObjectKey{Name: req.Name, Namespace: ns}, existing) - shouldRecreate := errors.IsNotFound(getErr) || - (getErr == nil && existing.Status.Phase == slessv1alpha1.ServicePhaseFailed) - if shouldRecreate { - if getErr == nil { - _ = h.K8s.Delete(r.Context(), existing) + + // --- Сценарий A: объект помечен на удаление --- + // DeletionTimestamp ≠ zero означает что k8s принял DELETE и ждёт finalizer. + // Finalizer снимает service_controller.go асинхронно (обычно 1-5 сек). + // Мы не можем создать объект пока старый существует — ждём его исчезновения. + if getErr == nil && !existing.DeletionTimestamp.IsZero() { + + deleted := false + // Проверяем каждую секунду до 30 итераций (= 30 секунд максимум). + // Обычно finalizer снимается за 1-3 секунды, 30 — запас на перегруженный кластер. + for i := 0; i < 30; i++ { + time.Sleep(1 * time.Second) + checkErr := h.K8s.Get(r.Context(), client.ObjectKey{Name: req.Name, Namespace: ns}, existing) + // IsNotFound = объект полностью исчез из etcd — можно создавать. + if errors.IsNotFound(checkErr) { + deleted = true + break + } + // Любая другая ошибка Get — продолжаем ждать (может быть временный сбой API). } + + // 30 секунд прошло, объект всё ещё не удалён. + // Вероятно контроллер завис или finalizer не снимается. + // Возвращаем 409 с понятным сообщением — клиент (terraform) должен retry. + if !deleted { + writeJSON(w, http.StatusConflict, errResp("service is being deleted, try again later")) + return + } + + // Объект исчез. Сбрасываем ResourceVersion чтобы k8s воспринял как новый объект. + // ResourceVersion заполняется при первом Get — при Create должен быть пустым. svc.ResourceVersion = "" if createErr := h.K8s.Create(r.Context(), svc); createErr != nil { writeJSON(w, http.StatusInternalServerError, errResp(createErr.Error())) @@ -154,11 +206,39 @@ func (h *Handler) CreateService(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusCreated, svcToResponse(svc)) return } + + // --- Сценарии B и C: split-brain кеша или объект в фазе Failed --- + // shouldRecreate=true если: + // - getErr = NotFound (кеш врёт — объекта в etcd нет) + // - объект есть, но в фазе Failed (предыдущий build провалился) + shouldRecreate := errors.IsNotFound(getErr) || + (getErr == nil && existing.Status.Phase == slessv1alpha1.ServicePhaseFailed) + if shouldRecreate { + // Если объект реально существует (Failed) — удаляем его перед пересозданием. + // Ошибку удаления игнорируем: даже если Delete упал, + // следующий Create либо пройдёт (объект исчез) либо снова упадёт с понятной ошибкой. + if getErr == nil { + _ = h.K8s.Delete(r.Context(), existing) + } + // Сбрасываем ResourceVersion — при split-brain etcd считает объект новым. + svc.ResourceVersion = "" + if createErr := h.K8s.Create(r.Context(), svc); createErr != nil { + writeJSON(w, http.StatusInternalServerError, errResp(createErr.Error())) + return + } + writeJSON(w, http.StatusCreated, svcToResponse(svc)) + return + } + + // Объект существует, DeletionTimestamp=zero, фаза не Failed. + // Это легитимный конфликт — сервис реально работает, клиент пытается создать дубликат. writeJSON(w, http.StatusConflict, errResp("service already exists")) return } - // k8s возвращает StatusInvalid (422) при нарушении enum-валидации CRD (например, неизвестный runtime) - // — маппим это в 400, а не в 500 + + // k8s возвращает StatusInvalid (422) при нарушении enum-валидации CRD. + // Например: runtime="java8" которого нет в openapi enum сервисного CRD. + // Маппим в 400 (Bad Request) а не в 500 — это ошибка клиента, не сервера. if errors.IsInvalid(err) { writeJSON(w, http.StatusBadRequest, errResp("invalid service spec: "+err.Error())) return diff --git a/internal/builder/builder.go b/internal/builder/builder.go index 5269a62..aea9c07 100644 --- a/internal/builder/builder.go +++ b/internal/builder/builder.go @@ -1,16 +1,22 @@ -// Изменено: 2026-03-11 +// Изменено: 2026-03-22 // Builder — запускает kaniko Job в k8s для сборки Docker образа из кода функции. // Почему kaniko, а не docker-in-docker (DinD): // kaniko не требует privileged контейнер, что безопаснее в managed кластере. // kaniko читает контекст сборки из S3 напрямую. // Workflow: FunctionController вызывает Build → EnsureProject (Harbor) → Job запускается → образ пушится в registry. +// ImageExists: перед сборкой проверяем наличие образа по тегу в registry. +// Если образ с таким тегом (= hash кода) уже существует — сборка пропускается. +// Это даёт идемпотентность destroy+apply: если код не менялся, deploy мгновенный. package builder import ( "context" "crypto/sha256" + "encoding/json" "fmt" + "net/http" + "strings" "time" batchv1 "k8s.io/api/batch/v1" @@ -33,12 +39,12 @@ type Builder struct { registryHost string // куда пушим образ (DockerHub: "naeel"; Harbor: "host") registryProject string // проект/org внутри registry (Harbor: project; DockerHub: пусто) registrySecret string // имя k8s Secret с docker-кредами для kaniko - s3Endpoint string // откуда kaniko берёт код - s3AccessKey string - s3SecretKey string - s3Bucket string - namespace string // namespace где запускаем build Job'ы - harborClient Projecter // nil если Harbor не используется + s3Endpoint string // откуда kaniko берёт код + s3AccessKey string + s3SecretKey string + s3Bucket string + namespace string // namespace где запускаем build Job'ы + harborClient Projecter // nil если Harbor не используется } // Config — параметры для создания Builder'а. @@ -47,12 +53,12 @@ type Config struct { RegistryHost string RegistryProject string // пусто = DockerHub-режим (2 уровня); задан = project-режим (3 уровня) RegistrySecret string // имя k8s Secret с .dockerconfigjson для пуша образов - S3Endpoint string - S3AccessKey string - S3SecretKey string - S3Bucket string - Namespace string - HarborClient Projecter // nil — Harbor не используется, EnsureProject пропускается + S3Endpoint string + S3AccessKey string + S3SecretKey string + S3Bucket string + Namespace string + HarborClient Projecter // nil — Harbor не используется, EnsureProject пропускается } // New создаёт новый Builder. @@ -63,12 +69,12 @@ func New(c client.Client, cfg Config) *Builder { registryHost: cfg.RegistryHost, registryProject: cfg.RegistryProject, registrySecret: cfg.RegistrySecret, - s3Endpoint: cfg.S3Endpoint, - s3AccessKey: cfg.S3AccessKey, - s3SecretKey: cfg.S3SecretKey, - s3Bucket: cfg.S3Bucket, - namespace: cfg.Namespace, - harborClient: cfg.HarborClient, + s3Endpoint: cfg.S3Endpoint, + s3AccessKey: cfg.S3AccessKey, + s3SecretKey: cfg.S3SecretKey, + s3Bucket: cfg.S3Bucket, + namespace: cfg.Namespace, + harborClient: cfg.HarborClient, } } @@ -103,6 +109,88 @@ func (b *Builder) ImageRef(namespace, funcName, s3Key string) string { return fmt.Sprintf("%s/%s-%s:%s", b.registryHost, nsPrefix, funcName, tag) } +// ImageExists проверяет наличие образа в Docker Registry v2 по тегу (без pull). +// Использует анонимный bearer-token для публичных репо (DockerHub). +// Возвращает true если образ с таким тегом уже запушен — сборка не нужна. +// +// Почему анонимный токен, а не credentials: +// +// DockerHub выдаёт pull-token без авторизации для публичных репо через +// GET /token?service=registry.docker.io&scope=repository:{repo}:pull +// Это стандартный Docker Registry v2 auth flow (RFC 7235). +func (b *Builder) ImageExists(ctx context.Context, imageRef string) bool { + // imageRef вида: "naeel/slessffd1-pg-search:47cab27ada70" + // или "host/project/func:tag" — разбираем по последнему ":" + colonIdx := strings.LastIndex(imageRef, ":") + if colonIdx < 0 { + return false + } + repoFull := imageRef[:colonIdx] + tag := imageRef[colonIdx+1:] + + // Определяем registry host и repo path. + // DockerHub: "naeel/sless-ff-pg" → registry = index.docker.io, repo = "naeel/sless-ff-pg" + // Приватный: "harbor.host/proj/func" → registry = "harbor.host", repo = "proj/func" + registryHost := "index.docker.io" + repo := repoFull + if parts := strings.SplitN(repoFull, "/", 3); len(parts) == 3 { + // host/project/name — кастомный registry + registryHost = parts[0] + repo = parts[1] + "/" + parts[2] + } + + // Получаем анонимный/публичный bearer-token для pull доступа к репо. + // DockerHub: https://auth.docker.io/token?service=registry.docker.io&scope=repository:{repo}:pull + // Для приватных registry этот шаг вернёт 401 → ImageExists вернёт false → пойдём строить. + tokenURL := fmt.Sprintf("https://auth.docker.io/token?service=registry.docker.io&scope=repository:%s:pull", repo) + if registryHost != "index.docker.io" { + // Для не-DockerHub registry: пробуем без токена (Harbor с allow anon push) + // Если 401 — false, пусть builder разберётся. + tokenURL = "" + } + + httpClient := &http.Client{Timeout: 5 * time.Second} + + var bearerToken string + if tokenURL != "" { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, tokenURL, nil) + if err != nil { + return false + } + resp, err := httpClient.Do(req) + if err != nil || resp.StatusCode != http.StatusOK { + return false + } + defer resp.Body.Close() + var tkResp struct { + Token string `json:"token"` + } + if err := json.NewDecoder(resp.Body).Decode(&tkResp); err != nil { + return false + } + bearerToken = tkResp.Token + } + + // HEAD /v2/{repo}/manifests/{tag} — проверяем наличие тега без скачивания слоёв. + manifestURL := fmt.Sprintf("https://%s/v2/%s/manifests/%s", registryHost, repo, tag) + req, err := http.NewRequestWithContext(ctx, http.MethodHead, manifestURL, nil) + if err != nil { + return false + } + // Docker Registry v2 требует Accept header для манифестов. + req.Header.Set("Accept", "application/vnd.docker.distribution.manifest.v2+json") + if bearerToken != "" { + req.Header.Set("Authorization", "Bearer "+bearerToken) + } + + resp, err := httpClient.Do(req) + if err != nil { + return false + } + defer resp.Body.Close() + return resp.StatusCode == http.StatusOK +} + // Build запускает kaniko Job для сборки образа функции. // Возвращает имя Job'а чтобы контроллер мог следить за его статусом. func (b *Builder) Build(ctx context.Context, namespace, funcName, s3Key string) (string, error) { diff --git a/main.go b/main.go index 35c4a23..61f104a 100644 --- a/main.go +++ b/main.go @@ -135,12 +135,12 @@ func main() { RegistryHost: cfg.RegistryHost, RegistryProject: cfg.RegistryProject, RegistrySecret: cfg.RegistrySecret, - S3Endpoint: cfg.S3Endpoint, - S3AccessKey: cfg.S3AccessKey, - S3SecretKey: cfg.S3SecretKey, - S3Bucket: cfg.S3Bucket, - Namespace: "sless", - HarborClient: harborProjecter, + S3Endpoint: cfg.S3Endpoint, + S3AccessKey: cfg.S3AccessKey, + S3SecretKey: cfg.S3SecretKey, + S3Bucket: cfg.S3Bucket, + Namespace: "sless", + HarborClient: harborProjecter, }) if err = (&controllers.FunctionReconciler{ diff --git a/runtimes/go1.23/Dockerfile b/runtimes/go1.23/Dockerfile index 9258a19..0ecac60 100644 --- a/runtimes/go1.23/Dockerfile +++ b/runtimes/go1.23/Dockerfile @@ -1,18 +1,20 @@ -# Изменено: 2026-03-19 — v0.1.1: добавлены go.sum + go mod download (pgx/v5 v5.7.2). +# Изменено: 2026-03-22 — v0.1.3: server/ субдиректория + go.work для поддержки пользовательских go.mod. # Base builder image для Go 1.23 serverless функций. -# Это BUILDER-образ: golang:1.23-alpine + server.go + pre-cached зависимости. +# Это BUILDER-образ: golang:1.23-alpine + server/ + pre-cached зависимости. # go mod download кеширует pgx/v5 в /root/go/pkg/mod — kaniko не скачивает их при каждой сборке. # # kaniko генерирует Dockerfile для каждой функции: -# FROM naeel/sless-runtime-go1.23:v0.1.1 AS builder +# FROM naeel/sless-runtime-go1.23:v0.1.3 AS builder # WORKDIR /app # COPY . /app/handler/ -# RUN CGO_ENABLED=0 go build -o /server . +# RUN [ -f /app/handler/go.mod ] || printf 'module sless/fn/handler\n\ngo 1.23\n' > /app/handler/go.mod +# RUN printf 'go 1.23\n\nuse ./server\nuse ./handler\n\nreplace sless/fn/handler => ./handler\n' > /app/go.work +# RUN CGO_ENABLED=0 go build -o /server ./server # FROM alpine:3.20 -# COPY --from=builder /server /server → финальный образ +# COPY --from=builder /server /server # -# Почему не multi-stage здесь: base образ должен иметь Go компилятор и module cache. -# Бинарник собирается kaniko, а не здесь. +# Почему server/ субдиректория: go.work + replace позволяет пользователю иметь любой go.mod +# с любыми зависимостями. Вложенные модули (nested) поддерживаются через workspace. FROM golang:1.23-alpine @@ -20,10 +22,10 @@ WORKDIR /app # Копируем mod файлы первыми — отдельный кешируемый слой. # При изменении go.mod/go.sum слой инвалидируется и go mod download запускается заново. -COPY go.mod go.sum /app/ -RUN go mod download +COPY server/go.mod server/go.sum /app/server/ +RUN cd /app/server && go mod download # server.go — HTTP-wrapper + job-runner (main package, импортирует sless/fn/handler). -COPY server.go /app/server.go +COPY server/server.go /app/server/server.go EXPOSE 8080 diff --git a/runtimes/go1.23/server/go.mod b/runtimes/go1.23/server/go.mod new file mode 100644 index 0000000..1470714 --- /dev/null +++ b/runtimes/go1.23/server/go.mod @@ -0,0 +1,15 @@ +// Изменено: 2026-03-22 — переименован в sless/fn/server для поддержки go.work + user go.mod +module sless/fn/server + +go 1.23 + +require github.com/jackc/pgx/v5 v5.7.2 + +require ( + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect + github.com/jackc/puddle/v2 v2.2.2 // indirect + golang.org/x/crypto v0.31.0 // indirect + golang.org/x/sync v0.10.0 // indirect + golang.org/x/text v0.21.0 // indirect +) diff --git a/runtimes/go1.23/server/go.sum b/runtimes/go1.23/server/go.sum new file mode 100644 index 0000000..731b5df --- /dev/null +++ b/runtimes/go1.23/server/go.sum @@ -0,0 +1,28 @@ +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +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/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.7.2 h1:mLoDLV6sonKlvjIEsV56SkWNCnuNv531l94GaIzO+XI= +github.com/jackc/pgx/v5 v5.7.2/go.mod h1:ncY89UGWxg82EykZUwSpUKEfccBGGYq1xjrOpsbsfGQ= +github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= +github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk= +github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= +golang.org/x/crypto v0.31.0 h1:ihbySMvVjLAeSH1IbfcRTkD/iNscyz8rGzjF/E5hV6U= +golang.org/x/crypto v0.31.0/go.mod h1:kDsLvtWBEx7MV9tJOj9bnXsPbxwJQ6csT/x4KIN4Ssk= +golang.org/x/sync v0.10.0 h1:3NQrjDixjgGwUOCaF8w2+VYHv0Ve/vGYSbdkTa98gmQ= +golang.org/x/sync v0.10.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= +golang.org/x/text v0.21.0 h1:zyQAAkrwaneQ066sspRyJaG9VNi/YJ1NfzcGB3hZ/qo= +golang.org/x/text v0.21.0/go.mod h1:4IBbMaMmOPCJ8SecivzSH54+73PCFmPWxNTLm+vZkEQ= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/runtimes/go1.23/server/server.go b/runtimes/go1.23/server/server.go new file mode 100644 index 0000000..71b64bd --- /dev/null +++ b/runtimes/go1.23/server/server.go @@ -0,0 +1,97 @@ +// Создано: 2026-03-11 +// Изменено: 2026-03-21 (v0.1.2) — panic recovery: при panic в Handle() возвращаем HTTP 500 вместо EOF. +// HTTP-обёртка для serverless функций на Go 1.23. +// Компилируется kaniko ВМЕСТЕ с пользовательским кодом (package handler). +// +// Интерфейс пользователя — файл handler.go: +// +// package handler +// +// func Handle(event map[string]interface{}) interface{} { +// return map[string]interface{}{"hello": event["name"]} +// } +// +// SLESS_MODE=job: разово вызвать Handle(event) → вывести JSON → выйти (для FunctionJob). +// По умолчанию: HTTP-сервер, каждый запрос → Handle(event) → JSON ответ. + +package main + +import ( + "encoding/json" + "fmt" + "io" + "log" + "net/http" + "os" + + "sless/fn/handler" +) + +func main() { + mode := os.Getenv("SLESS_MODE") + + // job-runner: разово вызвать Handle(event), вывести результат в JSON и выйти. + if mode == "job" { + eventJSON := os.Getenv("SLESS_EVENT") + if eventJSON == "" { + eventJSON = "{}" + } + var event map[string]interface{} + if err := json.Unmarshal([]byte(eventJSON), &event); err != nil { + event = map[string]interface{}{} + } + result := handler.Handle(event) + out, _ := json.Marshal(result) + fmt.Println(string(out)) + return + } + + // HTTP-сервер + port := "8080" + http.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, `{"status":"ok"}`) + }) + http.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + // Перехватываем panic в пользовательском коде — возвращаем 500 вместо EOF. + // Без recover() Go закрывает соединение при панике → ingress видит EOF → 502. + defer func() { + if rec := recover(); rec != nil { + log.Printf("panic recovered: %v", rec) + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusInternalServerError) + fmt.Fprintf(w, `{"error":"panic: %v"}`, rec) + } + }() + event := map[string]interface{}{} + if r.Method == http.MethodPost || r.Method == http.MethodPut { + body, err := io.ReadAll(r.Body) + if err == nil && len(body) > 0 { + _ = json.Unmarshal(body, &event) + } + } + event["_path"] = r.URL.Path + event["_method"] = r.Method + if q := r.URL.Query(); len(q) > 0 { + qmap := map[string]interface{}{} + for k, v := range q { + if len(v) == 1 { + qmap[k] = v[0] + } else { + qmap[k] = v + } + } + event["_query"] = qmap + } + result := handler.Handle(event) + out, err := json.Marshal(result) + if err != nil { + http.Error(w, `{"error":"marshal failed"}`, 500) + return + } + w.Header().Set("Content-Type", "application/json") + w.Write(out) + }) + log.Printf("sless runtime (go1.23) listening on :%s", port) + log.Fatal(http.ListenAndServe(":"+port, nil)) +}