feat(v0.1.58): ImageExists cache hit, timing analysis, in-cluster registry plan
This commit is contained in:
@@ -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))
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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)
|
||||
|
||||
### Симптом
|
||||
|
||||
+25
-1
@@ -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 за секунды.
|
||||
|
||||
---
|
||||
|
||||
|
||||
@@ -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 — тест устойчивости к сбоям.
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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),
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
}
|
||||
}
|
||||
@@ -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),
|
||||
}
|
||||
}
|
||||
@@ -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 ────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
// Сбрасываем ResourceVersion — при split-brain etcd считает объект новым
|
||||
// Любая другая ошибка (timeout, сбой API) — продолжаем ждать.
|
||||
}
|
||||
|
||||
// Таймаут: объект не исчез за 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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"
|
||||
@@ -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) {
|
||||
|
||||
+12
-10
@@ -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
|
||||
|
||||
@@ -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
|
||||
)
|
||||
@@ -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=
|
||||
@@ -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))
|
||||
}
|
||||
Reference in New Issue
Block a user