Compare commits
11
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
dc65f7ab8f | ||
|
|
df84eae75a | ||
|
|
45bd9389e5 | ||
|
|
d7fda15d35 | ||
|
|
8dc07445ac | ||
|
|
014b99ed16 | ||
|
|
d87981713d | ||
|
|
8a8b815492 | ||
|
|
0ebae25877 | ||
|
|
a379091b8a | ||
|
|
1aac3f5093 |
@@ -64,3 +64,4 @@ sless-plan
|
||||
plan.out
|
||||
sless-plan
|
||||
examples/.git
|
||||
event-dispatcher
|
||||
|
||||
@@ -0,0 +1,25 @@
|
||||
# Изменено: 2026-03-19
|
||||
# Multi-stage build для event-dispatcher.
|
||||
# Stage 1: сборка бинаря
|
||||
# Stage 2: минимальный образ (нужен ca-certificates для TLS к RabbitMQ и k8s API)
|
||||
FROM golang:1.25-alpine AS builder
|
||||
ARG TARGETOS
|
||||
ARG TARGETARCH
|
||||
|
||||
WORKDIR /workspace
|
||||
COPY go.mod go.mod
|
||||
COPY go.sum go.sum
|
||||
RUN go mod download
|
||||
|
||||
COPY api/ api/
|
||||
COPY services/event-dispatcher/ services/event-dispatcher/
|
||||
|
||||
RUN CGO_ENABLED=0 GOOS=${TARGETOS:-linux} GOARCH=${TARGETARCH} \
|
||||
go build -a -o event-dispatcher ./services/event-dispatcher/
|
||||
|
||||
FROM alpine:3.19
|
||||
RUN apk add --no-cache ca-certificates
|
||||
WORKDIR /
|
||||
COPY --from=builder /workspace/event-dispatcher .
|
||||
USER 65532:65532
|
||||
ENTRYPOINT ["/event-dispatcher"]
|
||||
@@ -1,6 +1,8 @@
|
||||
// Изменено: 2026-03-08
|
||||
// Описание CRD Trigger — триггер для функции (HTTP или Cron).
|
||||
// Изменено: 2026-03-19
|
||||
// Описание CRD Trigger — триггер для функции (HTTP, Cron или Event).
|
||||
// Один Trigger ссылается на одну Function и определяет способ вызова.
|
||||
// Event-тип: event-dispatcher подписывается на AMQP очередь и при сообщении
|
||||
// вызывает функцию по внутреннему HTTP.
|
||||
|
||||
package v1alpha1
|
||||
|
||||
@@ -16,6 +18,9 @@ const (
|
||||
TriggerTypeHTTP TriggerType = "http"
|
||||
// TriggerTypeCron — функция вызывается по расписанию (k8s CronJob)
|
||||
TriggerTypeCron TriggerType = "cron"
|
||||
// TriggerTypeEvent — функция вызывается при получении сообщения из AMQP очереди.
|
||||
// event-dispatcher подписывается на spec.queue в RabbitMQ и делает POST на HTTP endpoint функции.
|
||||
TriggerTypeEvent TriggerType = "event"
|
||||
)
|
||||
|
||||
// TriggerSpec — желаемое состояние триггера.
|
||||
@@ -30,8 +35,8 @@ type TriggerSpec struct {
|
||||
// +kubebuilder:validation:Required
|
||||
FunctionRef string `json:"functionRef"`
|
||||
|
||||
// Type — тип триггера: http или cron
|
||||
// +kubebuilder:validation:Enum=http;cron
|
||||
// Type — тип триггера: http, cron или event
|
||||
// +kubebuilder:validation:Enum=http;cron;event
|
||||
// +kubebuilder:validation:Required
|
||||
Type TriggerType `json:"type"`
|
||||
|
||||
@@ -42,6 +47,10 @@ type TriggerSpec struct {
|
||||
// Актуально для cron: запускаем pod заранее чтобы избежать cold start.
|
||||
// +kubebuilder:default=300
|
||||
PreWarmSeconds int32 `json:"preWarmSeconds,omitempty"`
|
||||
|
||||
// Queue — имя AMQP очереди в RabbitMQ (только для type=event).
|
||||
// event-dispatcher подпишется на эту очередь и вызовет функцию при каждом сообщении.
|
||||
Queue string `json:"queue,omitempty"`
|
||||
}
|
||||
|
||||
// TriggerStatus — наблюдаемое состояние триггера (заполняет контроллер).
|
||||
@@ -65,6 +74,7 @@ type TriggerStatus struct {
|
||||
//+kubebuilder:printcolumn:name="Function",type=string,JSONPath=`.spec.functionRef`
|
||||
//+kubebuilder:printcolumn:name="Active",type=boolean,JSONPath=`.status.active`
|
||||
//+kubebuilder:printcolumn:name="URL",type=string,JSONPath=`.status.url`
|
||||
//+kubebuilder:printcolumn:name="Queue",type=string,JSONPath=`.spec.queue`
|
||||
//+kubebuilder:printcolumn:name="Age",type=date,JSONPath=`.metadata.creationTimestamp`
|
||||
|
||||
// Trigger — ресурс для управления способом вызова Function.
|
||||
|
||||
@@ -124,6 +124,9 @@ func (r *TriggerReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ct
|
||||
case slessv1alpha1.TriggerTypeCron:
|
||||
logger.Info("reconcile cron trigger", "trigger", tr.Name)
|
||||
return r.reconcileCron(ctx, tr, fn)
|
||||
case slessv1alpha1.TriggerTypeEvent:
|
||||
logger.Info("reconcile event trigger", "trigger", tr.Name)
|
||||
return r.reconcileEvent(ctx, tr, fn)
|
||||
}
|
||||
|
||||
return ctrl.Result{}, nil
|
||||
@@ -325,3 +328,44 @@ func (r *TriggerReconciler) SetupWithManager(mgr ctrl.Manager) error {
|
||||
For(&slessv1alpha1.Trigger{}).
|
||||
Complete(r)
|
||||
}
|
||||
|
||||
// reconcileEvent обрабатывает Trigger{type:event}.
|
||||
// Оператор не управляет AMQP напрямую — это задача event-dispatcher.
|
||||
// Здесь: убеждаемся что Service функции существует (dispatcher использует его для POST),
|
||||
// обновляем статус триггера.
|
||||
func (r *TriggerReconciler) reconcileEvent(ctx context.Context, tr *slessv1alpha1.Trigger, fn *slessv1alpha1.Function) (ctrl.Result, error) {
|
||||
if tr.Spec.Queue == "" {
|
||||
tr.Status.Active = false
|
||||
tr.Status.Message = "queue is required for type=event"
|
||||
_ = r.Status().Update(ctx, tr)
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
deployNS := "sless-fn-" + tr.Namespace
|
||||
|
||||
// Service нужен event-dispatcher для доставки сообщений в функцию по HTTP.
|
||||
// Имя Service совпадает с именем Function — dispatcher строит URL как
|
||||
// http://{functionRef}.{deployNS}.svc.cluster.local:8080/
|
||||
wantSvc := &corev1.Service{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: fn.Name, Namespace: deployNS},
|
||||
Spec: corev1.ServiceSpec{
|
||||
Selector: map[string]string{"app": fn.Name},
|
||||
Ports: []corev1.ServicePort{{Port: 8080, Protocol: corev1.ProtocolTCP}},
|
||||
},
|
||||
}
|
||||
existingSvc := &corev1.Service{}
|
||||
if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, existingSvc); err != nil {
|
||||
if errors.IsNotFound(err) {
|
||||
if err := r.Create(ctx, wantSvc); err != nil {
|
||||
return ctrl.Result{}, fmt.Errorf("create service for event trigger: %w", err)
|
||||
}
|
||||
} else {
|
||||
return ctrl.Result{}, fmt.Errorf("get service: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
tr.Status.Active = true
|
||||
tr.Status.Message = fmt.Sprintf("listening on queue %q via event-dispatcher", tr.Spec.Queue)
|
||||
_ = r.Status().Update(ctx, tr)
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
# Изменено: 2026-03-19
|
||||
# event-dispatcher — отдельный сервис для обработки event-триггеров.
|
||||
# Следит за Trigger CRD{type:event}, подписывается на AMQP очереди,
|
||||
# при сообщении делает POST на внутренний HTTP endpoint функции.
|
||||
#
|
||||
# Требует:
|
||||
# - SecretRef: sless-operator-secret (RABBITMQ_URL)
|
||||
# - ClusterRole: event-dispatcher-role (чтение Trigger CRD, Namespace)
|
||||
---
|
||||
apiVersion: v1
|
||||
kind: ServiceAccount
|
||||
metadata:
|
||||
name: event-dispatcher
|
||||
namespace: sless
|
||||
---
|
||||
apiVersion: rbac.authorization.k8s.io/v1
|
||||
kind: ClusterRole
|
||||
metadata:
|
||||
name: event-dispatcher-role
|
||||
rules:
|
||||
# Нужно читать Trigger CRD по всем namespace (event-dispatcher глобальный)
|
||||
- apiGroups: ["sless.kube5s.ru"]
|
||||
resources: ["triggers"]
|
||||
verbs: ["get", "list", "watch"]
|
||||
# Нужно читать namespace для построения URLs функций
|
||||
- apiGroups: [""]
|
||||
resources: ["namespaces"]
|
||||
verbs: ["get", "list", "watch"]
|
||||
---
|
||||
apiVersion: rbac.authorization.k8s.io/v1
|
||||
kind: ClusterRoleBinding
|
||||
metadata:
|
||||
name: event-dispatcher-rolebinding
|
||||
subjects:
|
||||
- kind: ServiceAccount
|
||||
name: event-dispatcher
|
||||
namespace: sless
|
||||
roleRef:
|
||||
kind: ClusterRole
|
||||
name: event-dispatcher-role
|
||||
apiGroup: rbac.authorization.k8s.io
|
||||
---
|
||||
apiVersion: apps/v1
|
||||
kind: Deployment
|
||||
metadata:
|
||||
name: event-dispatcher
|
||||
namespace: sless
|
||||
labels:
|
||||
app: event-dispatcher
|
||||
spec:
|
||||
replicas: 1
|
||||
selector:
|
||||
matchLabels:
|
||||
app: event-dispatcher
|
||||
template:
|
||||
metadata:
|
||||
labels:
|
||||
app: event-dispatcher
|
||||
spec:
|
||||
serviceAccountName: event-dispatcher
|
||||
containers:
|
||||
- name: event-dispatcher
|
||||
image: naeel/sless-event-dispatcher:v0.1.0
|
||||
imagePullPolicy: Always
|
||||
env:
|
||||
- name: RABBITMQ_URL
|
||||
valueFrom:
|
||||
secretKeyRef:
|
||||
name: sless-operator-secret
|
||||
key: RABBITMQ_URL
|
||||
resources:
|
||||
requests:
|
||||
cpu: 50m
|
||||
memory: 64Mi
|
||||
limits:
|
||||
cpu: 200m
|
||||
memory: 128Mi
|
||||
@@ -130,6 +130,8 @@ metadata:
|
||||
cert-manager.io/cluster-issuer: letsencrypt-prod
|
||||
nginx.ingress.kubernetes.io/force-ssl-redirect: "true"
|
||||
nginx.ingress.kubernetes.io/ssl-redirect: "true"
|
||||
nginx.ingress.kubernetes.io/proxy-read-timeout: "900"
|
||||
nginx.ingress.kubernetes.io/proxy-send-timeout: "900"
|
||||
spec:
|
||||
ingressClassName: nginx
|
||||
rules:
|
||||
|
||||
@@ -9,7 +9,7 @@
|
||||
- **Репозиторий:** `gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless`
|
||||
- **Локальная копия:** `/home/naeel/remote_dev/sless/`
|
||||
- **Remote server:** `naeel@5.172.178.213` (workspace: `~/terra/sless/`)
|
||||
- **SSH ключ:** `/home/naeel/remote_dev/sless/secrets/naeel_vm_id_ed25519`
|
||||
- **SSH ключ:** `/home/naeel/.ssh/naeel_vm_id_ed25519`
|
||||
- **SSH команда:** `ssh -i <ключ> -o StrictHostKeyChecking=no naeel@5.172.178.213`
|
||||
- **Активная ветка:** `feat/web-console` (последний коммит `a04dfb2`)
|
||||
- **Git origin:** `https://gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless.git`
|
||||
@@ -320,8 +320,8 @@ kubectl set env deployment/sless-funcs-service -n sless SLESS_SERVICE_TOKEN=<jwt
|
||||
## 13. Деплой нового образа — стандартный workflow
|
||||
|
||||
```bash
|
||||
SSH="ssh -i /home/naeel/remote_dev/sless/secrets/naeel_vm_id_ed25519 -o StrictHostKeyChecking=no naeel@5.172.178.213"
|
||||
SCP="scp -i /home/naeel/remote_dev/sless/secrets/naeel_vm_id_ed25519 -o StrictHostKeyChecking=no"
|
||||
SSH="ssh -i /home/naeel/.ssh/naeel_vm_id_ed25519 -o StrictHostKeyChecking=no naeel@5.172.178.213"
|
||||
SCP="scp -i /home/naeel/.ssh/naeel_vm_id_ed25519 -o StrictHostKeyChecking=no"
|
||||
|
||||
# 1. Скопировать изменённые файлы оператора на remote
|
||||
$SCP /home/naeel/remote_dev/sless/path/to/file.go naeel@5.172.178.213:~/terra/sless/path/to/file.go
|
||||
|
||||
+170
-2
@@ -1,5 +1,104 @@
|
||||
# Решения и обоснования
|
||||
|
||||
---
|
||||
|
||||
## 2026-03-19 — Go runtime v0.1.1: внешние зависимости через go.mod/go.sum
|
||||
|
||||
### Контекст
|
||||
|
||||
Go runtime `naeel/sless-runtime-go1.23:v0.1.0` содержал `go.mod` только с `module sless/fn` и `go 1.23`.
|
||||
Никаких `require` — пользовательский код мог использовать только stdlib.
|
||||
|
||||
При попытке добавить `pgxpool` в handler.go функция не собиралась (зависимость не найдена).
|
||||
|
||||
### Решение
|
||||
|
||||
Добавить `require github.com/jackc/pgx/v5 v5.7.2` в `runtimes/go1.23/go.mod`.
|
||||
Сгенерировать `go.sum` через `go mod tidy` (stub `.go` файл с импортом нужен — иначе tidy удалит deps).
|
||||
Обновить `Dockerfile` — добавить `COPY go.sum` + `RUN go mod download` **до** копирования пользовательского кода → зависимости кешируются в слое Docker, не скачиваются при каждой сборке функции.
|
||||
|
||||
### Почему pgx/v5, а не lib/pq
|
||||
|
||||
- `pgx/v5` — современный нативный PG-драйвер, `pgxpool` встроен, не нужен отдельный `database/sql`
|
||||
- `lib/pq` — legacy, минимальный API, отсутствует connection pool
|
||||
- `jackc/pgx/v5 v5.7.2` — последний стабильный тег на момент решения
|
||||
|
||||
### Что стало возможным
|
||||
|
||||
Любая Go функция в платформе может импортировать `pgxpool` и работать с PG напрямую:
|
||||
```go
|
||||
import "github.com/jackc/pgx/v5/pgxpool"
|
||||
```
|
||||
|
||||
### Версионирование образа
|
||||
|
||||
`v0.1.0` → `v0.1.1` — изменение breaking: бинарник пересобирается с новыми deps.
|
||||
Base image в `context.go` обновляется с `v0.1.0` на `v0.1.1`, оператор бампится.
|
||||
|
||||
---
|
||||
|
||||
## 2026-03-19 — Архитектура event-trigger (Вариант A: отдельный event-dispatcher)
|
||||
|
||||
### Контекст
|
||||
|
||||
До этого event-monitor/writer/cleaner работали как пользовательские sless-функции
|
||||
в namespace юзера — это неправильно: они ходили в операторскую Postgres напрямую,
|
||||
создавали таблицы без миграций, зависели от self-hosted rabbitmq.
|
||||
Всё это удалено из кластера (audit 2026-03-19).
|
||||
|
||||
### Варианты которые рассматривались
|
||||
|
||||
**Вариант A: отдельный event-dispatcher сервис** ← ВЫБРАН
|
||||
**Вариант B: dispatcher встроен горутиной в оператор**
|
||||
**Вариант C: CronJob polling из очереди**
|
||||
|
||||
### Решение: Вариант A
|
||||
|
||||
**Почему не B:** AMQP-соединения внутри operator reconciler усложняют lifecycle
|
||||
и тестирование. Падение AMQP затронет весь оператор.
|
||||
|
||||
**Почему не C:** polling — не realtime, не масштабируется, неловкий ACK.
|
||||
|
||||
**Почему A:** чистое разделение ответственности. Оператор управляет CRD,
|
||||
dispatcher управляет AMQP. Независимые restart/deploy. Легко тестировать отдельно.
|
||||
|
||||
### Поток данных
|
||||
|
||||
```
|
||||
Пользователь:
|
||||
kubectl apply — Trigger{type:event, queue:"orders", functionRef:"my-func"}
|
||||
|
||||
sless-operator (trigger_controller.go):
|
||||
reconcileEvent → валидирует что Function существует
|
||||
→ устанавливает status.active = true
|
||||
|
||||
event-dispatcher (services/event-dispatcher/):
|
||||
k8s informer наблюдает Trigger CRD по всем namespace
|
||||
При type=event → amqp.Channel.Consume(spec.queue)
|
||||
При сообщении → POST http://<fn-svc>.<fn-ns>.svc.cluster.local:8080/
|
||||
→ 2xx → ack
|
||||
→ не 2xx / timeout → nack (requeue)
|
||||
При удалении Trigger → закрыть consumer
|
||||
```
|
||||
|
||||
### Что меняется в коде
|
||||
|
||||
| Файл | Изменение |
|
||||
|------|-----------|
|
||||
| `api/v1alpha1/trigger_types.go` | +TriggerTypeEvent, +Queue в TriggerSpec |
|
||||
| `controllers/trigger_controller.go` | +reconcileEvent (валидация + status) |
|
||||
| `internal/config/config.go` | +RabbitMQURL |
|
||||
| `services/event-dispatcher/` | новый Go-сервис (main + dispatcher + watcher) |
|
||||
| `deployments/k8s/event-dispatcher.yaml` | Deployment + ServiceAccount + ClusterRole |
|
||||
|
||||
### Инфраструктура
|
||||
|
||||
RabbitMQ: managed через Nubes (Вариант A требует стабильного брокера).
|
||||
- Управляется rabbitmq-operator в namespace `operators`
|
||||
- namespace: `1dbfe9da-ce1c-4958-b359-d016a4b455c8`
|
||||
- host: `rabbitmqk8s.1dbfe9da-ce1c-4958-b359-d016a4b455c8.svc.cluster.local`
|
||||
- credentials: в `sless-operator-secret` (RABBITMQ_URL) — добавить при деплое
|
||||
|
||||
## 2026-03-18 — Архитектура: funcs как глобальный сервис, web-консоль
|
||||
|
||||
### Хранение кода функций
|
||||
@@ -105,12 +204,12 @@ provider "sless" {
|
||||
**Как запускать:**
|
||||
```bash
|
||||
# Сначала синхронизировать изменения:
|
||||
rsync -av -e "ssh -i /home/naeel/remote_dev/common/id_ed25519.txt" \
|
||||
rsync -av -e "ssh -i /home/naeel/.ssh/naeel_vm_id_ed25519" \
|
||||
/home/naeel/remote_dev/sless/examples/<example>/ \
|
||||
naeel@5.172.178.213:/home/naeel/terra/sless/examples/<example>/
|
||||
|
||||
# Затем запускать на remote:
|
||||
ssh -i /home/naeel/remote_dev/common/id_ed25519.txt naeel@5.172.178.213 \
|
||||
ssh -i /home/naeel/.ssh/naeel_vm_id_ed25519 naeel@5.172.178.213 \
|
||||
'cd /home/naeel/terra/sless/examples/<example> && terraform apply -auto-approve -no-color'
|
||||
```
|
||||
|
||||
@@ -762,3 +861,72 @@ for _, k := range keys { envVars = append(envVars, corev1.EnvVar{Name: k, Value:
|
||||
|
||||
**Тесты:** 2 теста в `controllers/function_controller_unit_test.go`
|
||||
(4 env vars → алфавитный порядок после SLESS_ENTRYPOINT; пустой Env → только SLESS_ENTRYPOINT).
|
||||
|
||||
---
|
||||
|
||||
## 2026-03-19 — pgx/v5 как PG-драйвер для Go функций (vs database/sql + lib/pq)
|
||||
|
||||
**Контекст:** Go runtime v0.1.1 — добавляем прямой доступ к PostgreSQL из функций.
|
||||
Нужно выбрать: database/sql + lib/pq, или чистый pgx/v5?
|
||||
|
||||
**Решение:** Использовать `github.com/jackc/pgx/v5` напрямую, без обёртки database/sql.
|
||||
|
||||
**Причины:**
|
||||
|
||||
1. **pgxpool из коробки** — `pgxpool.New()` без дополнительных пакетов. lib/pq требует `sql.Open` + настройку пула через `db.SetMaxOpenConns` и т.д.
|
||||
|
||||
2. **Нативный протокол PostgreSQL** — pgx реализует wire protocol напрямую, без CGO.
|
||||
lib/pq тоже pure Go, но pgx быстрее (~20% в бенчмарках) и активнее поддерживается.
|
||||
|
||||
3. **Контекст-нативность** — `pgxpool.Pool.Query(ctx, ...)` — context как первый аргумент везде.
|
||||
В database/sql контекст пришёл только в Go 1.8 как `QueryContext` — неудобный retrofit.
|
||||
|
||||
4. **Сканирование строк** — `pgx.CollectRows`, `pgx.ForEachRow` — удобнее чем `rows.Scan`.
|
||||
|
||||
5. **Экосистема** — pgx — де-факто стандарт в Go+PG проектах (используется в pgx, pgvector, ent).
|
||||
|
||||
**Что добавлено в рантайм:**
|
||||
```
|
||||
github.com/jackc/pgx/v5 v5.7.2
|
||||
github.com/jackc/pgpassfile v1.0.0 // indirect
|
||||
github.com/jackc/pgservicefile v0.0.0-... // indirect
|
||||
github.com/jackc/puddle/v2 v2.2.2 // indirect (connection pool)
|
||||
golang.org/x/crypto v0.31.0 // indirect (scram auth)
|
||||
golang.org/x/sync v0.10.0 // indirect
|
||||
golang.org/x/text v0.21.0 // indirect
|
||||
```
|
||||
|
||||
**go mod download** добавлен в Dockerfile до COPY server.go — слой с зависимостями кешируется отдельно.
|
||||
Пересборка функции (только изменение handler.go) не перекачивает ~15MB зависимостей.
|
||||
|
||||
---
|
||||
|
||||
## 2026-03-19 — Динамический таймаут в invoke.go из Function.Spec.TimeoutSec
|
||||
|
||||
**Контекст:** invoke.go проксирует HTTP-запросы к подам функций. До этого — глобальный `http.Client{Timeout: 30s}`.
|
||||
|
||||
**Проблема:** 30s — константа времени написания кода. Функции с `timeout_sec=700` (stress-тесты, batch-задачи) падают с `context deadline exceeded` раньше чем успевают завершиться.
|
||||
|
||||
**Решение:** Перед каждым вызовом читать `Function.Spec.TimeoutSec` из k8s и создавать `http.Client` с таймаутом = `TimeoutSec + 5s`.
|
||||
|
||||
**Почему +5s буфер:**
|
||||
- Нельзя ставить ровно `TimeoutSec` — есть сетевые задержки, TLS handshake, время на DNS резолв внутри кластера.
|
||||
- 5s достаточно для любых сетевых задержек в локальном k8s кластере.
|
||||
- Если функция реально завис на TimeoutSec — runtime сам должен прервать работу (это ответственность функции, не прокси).
|
||||
|
||||
**Почему не кешировать http.Client:**
|
||||
- Каждый вызов может прийти к разной функции с разным TimeoutSec.
|
||||
- http.Client создаётся дёшево — только структура с одним полем Timeout.
|
||||
- Кеш потребовал бы sync.Map или mutex — лишняя сложность без измеримой пользы.
|
||||
|
||||
**Деградация при недоступности k8s:**
|
||||
```go
|
||||
if err := h.K8s.Get(r.Context(), client.ObjectKey{...}, fn); err == nil {
|
||||
timeoutSec = fn.Spec.TimeoutSec
|
||||
}
|
||||
// если Get упал — timeoutSec=0 → invokeHTTPClient вернёт 30s (дефолт)
|
||||
```
|
||||
Это осознанный выбор: если мы не можем прочитать функцию — мы не знаем её таймаут,
|
||||
используем разумный дефолт вместо возврата ошибки.
|
||||
|
||||
**Коммит:** `d7fda15`, оператор `v0.1.40`
|
||||
|
||||
+125
-1
@@ -97,7 +97,7 @@ provider "sless" {
|
||||
|
||||
**Правило:** Все `terraform init/plan/apply/destroy` для `examples/` — только через SSH на `naeel@5.172.178.213`:
|
||||
```bash
|
||||
ssh -i /home/naeel/remote_dev/common/id_ed25519.txt naeel@5.172.178.213 \
|
||||
ssh -i /home/naeel/.ssh/naeel_vm_id_ed25519 naeel@5.172.178.213 \
|
||||
'cd /home/naeel/terra/sless/examples/<example> && terraform apply -auto-approve -no-color'
|
||||
```
|
||||
|
||||
@@ -631,3 +631,127 @@ Attribute runtime value must be one of: ["nodejs20" "python3.11" "go1.21"], got:
|
||||
### Версии
|
||||
- Оператор: `naeel/sless-operator:v0.1.13`
|
||||
- Провайдер: `terra.k8c.ru/naeel/sless v0.1.7`
|
||||
|
||||
---
|
||||
|
||||
## 2026-03-19 — Баг 4: Хардкодный 30s таймаут в invoke.go → context deadline exceeded
|
||||
|
||||
### Симптом
|
||||
|
||||
Вызов функции `stress-go-pgstorm` с `duration_sec=30` возвращал:
|
||||
```json
|
||||
{"error": "function unreachable: ... context deadline exceeded"}
|
||||
```
|
||||
При этом под был `Running`, логи показывали нормальную работу pgxpool.
|
||||
|
||||
### Корневая причина
|
||||
|
||||
В `internal/api/handler/invoke.go` (строка 24) был глобальный http.Client:
|
||||
```go
|
||||
var httpClient = &http.Client{Timeout: 30 * time.Second}
|
||||
```
|
||||
Функция реально отрабатывала ровно 30 секунд (duration_sec=30) + накладные расходы
|
||||
на pgxpool.New() и первый коннект к БД ≈ 1-2 секунды.
|
||||
Итого запрос превышал 30s → оператор разрывал соединение раньше чем функция успевала ответить.
|
||||
|
||||
### Почему так было написано
|
||||
|
||||
При создании invoke.go в марте 2026 таймаут 30s считался "достаточным для холодного
|
||||
старта". Длительные функции тогда не планировались. Когда появились batch/stress задачи
|
||||
с timeout_sec=600-700 — баг стал критическим.
|
||||
|
||||
### Решение
|
||||
|
||||
Убрать глобальный `httpClient`. Перед каждым вызовом:
|
||||
1. Получить Function CRD из k8s: `h.K8s.Get(ctx, ObjectKey{name, ns}, fn)`
|
||||
2. Прочитать `fn.Spec.TimeoutSec`
|
||||
3. Создать `http.Client{Timeout: TimeoutSec*time.Second + 5*time.Second}`
|
||||
4. Если функция не найдена (Get вернул ошибку) — дефолт 30s
|
||||
|
||||
```go
|
||||
func invokeHTTPClient(timeoutSec int32) *http.Client {
|
||||
t := time.Duration(timeoutSec)*time.Second + 5*time.Second
|
||||
if timeoutSec <= 0 {
|
||||
t = 30 * time.Second
|
||||
}
|
||||
return &http.Client{Timeout: t}
|
||||
}
|
||||
```
|
||||
|
||||
### Файл
|
||||
|
||||
`internal/api/handler/invoke.go` — исправлено в коммите `d7fda15`
|
||||
Оператор пересобран: `naeel/sless-operator:v0.1.40`
|
||||
|
||||
### Урок
|
||||
|
||||
**Никогда не хардкодить таймауты** в прокси-слое. Таймаут всегда должен браться
|
||||
из конфигурации вызываемого ресурса. `Function.Spec.TimeoutSec` существует именно для этого.
|
||||
|
||||
---
|
||||
|
||||
## 2026-03-19 — Баг 5: nginx ingress proxy-read-timeout не задан → 504 Gateway Time-out
|
||||
|
||||
### Симптом
|
||||
|
||||
После фикса invoke.go (баг 4) — повторный вызов `stress-go-pgstorm` вернул:
|
||||
```html
|
||||
<html><head><title>504 Gateway Time-out</title></head>
|
||||
<body><center><h1>504 Gateway Time-out</h1></center>
|
||||
<hr><center>nginx</center></body></html>
|
||||
```
|
||||
curl exit code 5 (не 0), процесс завершился через ~60 секунд после старта запроса.
|
||||
|
||||
### Корневая причина
|
||||
|
||||
У ingress `sless-operator` не было аннотации `proxy-read-timeout`.
|
||||
nginx ingress controller использует дефолт **60 секунд** если аннотация отсутствует.
|
||||
|
||||
```yaml
|
||||
# Было — аннотаций timeout нет вообще:
|
||||
annotations:
|
||||
kubernetes.io/ingress.class: nginx
|
||||
nginx.ingress.kubernetes.io/force-ssl-redirect: "true"
|
||||
nginx.ingress.kubernetes.io/ssl-redirect: "true"
|
||||
```
|
||||
|
||||
Цепочка: клиент → nginx (60s timeout) → оператор (705s) → функция (600s).
|
||||
Nginx оборвал соединение на 60-й секунде, хотя и оператор и функция были живы.
|
||||
|
||||
### Диагностика
|
||||
|
||||
Проверили аннотации всех ingress в ns sless:
|
||||
- `nodered` ingress: `proxy-read-timeout: "3600"` ✅ (кто-то правильно настроил)
|
||||
- `sless-funcs-ingress`: нет timeout аннотаций
|
||||
- `sless-operator`: нет timeout аннотаций ← **виновник**
|
||||
|
||||
### Решение
|
||||
|
||||
1. `kubectl annotate` для мгновенного применения:
|
||||
```bash
|
||||
kubectl annotate ingress sless-operator -n sless \
|
||||
nginx.ingress.kubernetes.io/proxy-read-timeout="900" \
|
||||
nginx.ingress.kubernetes.io/proxy-send-timeout="900" --overwrite
|
||||
```
|
||||
|
||||
2. Сохранить в манифест `deployments/k8s/operator.yaml`:
|
||||
```yaml
|
||||
nginx.ingress.kubernetes.io/proxy-read-timeout: "900"
|
||||
nginx.ingress.kubernetes.io/proxy-send-timeout: "900"
|
||||
```
|
||||
|
||||
### Почему 900s
|
||||
|
||||
- function timeout_sec = 700 → оператор ждёт 705s
|
||||
- nginx должен ждать дольше чем оператор → 900s с запасом
|
||||
- Не ставим 3600s как у nodered — избыточно для функций
|
||||
|
||||
### Файл
|
||||
|
||||
`deployments/k8s/operator.yaml` — обновлено в коммите `d7fda15`
|
||||
|
||||
### Урок
|
||||
|
||||
При развёртывании нового ingress **всегда явно задавать** `proxy-read-timeout`
|
||||
и `proxy-send-timeout`. Nginx дефолт 60s подходит только для быстрых API.
|
||||
Для любых операций дольше 30s — обязательны явные таймауты.
|
||||
|
||||
+2
-2
@@ -104,12 +104,12 @@ provider "sless" {
|
||||
**Как запускать:**
|
||||
```bash
|
||||
# Сначала синхронизировать изменения:
|
||||
rsync -av -e "ssh -i /home/naeel/remote_dev/common/id_ed25519.txt" \
|
||||
rsync -av -e "ssh -i /home/naeel/.ssh/naeel_vm_id_ed25519" \
|
||||
/home/naeel/remote_dev/sless/examples/<example>/ \
|
||||
naeel@5.172.178.213:/home/naeel/terra/sless/examples/<example>/
|
||||
|
||||
# Затем запускать на remote:
|
||||
ssh -i /home/naeel/remote_dev/common/id_ed25519.txt naeel@5.172.178.213 \
|
||||
ssh -i /home/naeel/.ssh/naeel_vm_id_ed25519 naeel@5.172.178.213 \
|
||||
'cd /home/naeel/terra/sless/examples/<example> && terraform apply -auto-approve -no-color'
|
||||
```
|
||||
|
||||
|
||||
+123
-1
@@ -1,6 +1,128 @@
|
||||
# Прогресс разработки
|
||||
|
||||
Последнее обновление: 2026-03-19 09:00
|
||||
Последнее обновление: 2026-03-19 22:30
|
||||
|
||||
---
|
||||
|
||||
## 2026-03-19 — Go runtime v0.1.1: pgx/v5 + stress-go-pgstorm ✅ ЗАВЕРШЕНО
|
||||
|
||||
### Цель
|
||||
Новая функция `stress-go-pgstorm` — Go код со 100 горутинами, которые 10 минут
|
||||
долбят PostgreSQL напрямую через `pgxpool`. Проверяем: Go runtime под нагрузкой,
|
||||
connection pool под конкурентными запросами, устойчивость кластера.
|
||||
|
||||
**Попутно выявлены и исправлены два системных бага:**
|
||||
- Хардкодный таймаут 30s в invoke.go (прокси оператора)
|
||||
- Отсутствие proxy-read-timeout на nginx ingress sless-operator
|
||||
|
||||
### Задачи
|
||||
|
||||
| # | Задача | Статус | Заметки |
|
||||
|---|--------|--------|---------|
|
||||
| 1 | `runtimes/go1.23/go.mod` — добавить `pgx/v5 v5.7.2` | ✅ | go mod tidy на remote, go.sum сгенерирован (28 строк) |
|
||||
| 2 | `runtimes/go1.23/Dockerfile` — `go mod download` кешируем deps | ✅ | Один stage `golang:1.23-alpine`, COPY go.mod+go.sum перед server.go |
|
||||
| 3 | Собрать `naeel/sless-runtime-go1.23:v0.1.1`, запушить | ✅ | docker build + push, sha256 verified |
|
||||
| 4 | `internal/builder/context.go` — тег `go1.23` v0.1.0 → v0.1.1 | ✅ | строка `return "naeel/sless-runtime-go1.23:v0.1.1", nil` |
|
||||
| 5 | `internal/api/handler/invoke.go` — динамический таймаут | ✅ | Таймаут = Function.Spec.TimeoutSec + 5s (был хардкод 30s) |
|
||||
| 6 | Пересобрать и задеплоить оператор `v0.1.40` | ✅ | rodlout complete |
|
||||
| 7 | `deployments/k8s/operator.yaml` — nginx proxy-read-timeout=900s | ✅ | kubectl annotate + манифест обновлён |
|
||||
| 8 | `code/stress-go-pgstorm/handler.go` — горутины + pgxpool | ✅ | 100 горутин, INSERT/COUNT/MAX, параметры workers/duration_sec/max_delay_ms |
|
||||
| 9 | `resources.tf` — новая функция + trigger | ✅ | timeout_sec=700, memory_mb=256, все PG env vars |
|
||||
| 10 | terraform apply — 2 ресурса добавлено | ✅ | apply complete: 2 added |
|
||||
| 11 | Запуск 10-минутного стресс-теста | ✅ | workers=100, duration_sec=600, err_ops=0, ok_ops=45918, ops_per_sec=76.5 — PostgreSQL выдержал |
|
||||
| 12 | Коммит d7fda15 | ✅ | feat: Go runtime v0.1.1 (pgx/v5)... |
|
||||
|
||||
### Архитектура функции stress-go-pgstorm
|
||||
|
||||
```
|
||||
Handle(event) → запускает N горутин (default 100)
|
||||
каждая горутина в цикле duration_sec (default 600):
|
||||
- случайная задержка 0-300ms (max_delay_ms)
|
||||
- чередующиеся операции: INSERT / SELECT COUNT / SELECT MAX
|
||||
- считает ok/err атомарно через sync/atomic
|
||||
WaitGroup.Wait() → возвращает итог
|
||||
возвращает: {runtime, version, workers, duration_sec, elapsed_sec,
|
||||
total_ops, ok_ops, err_ops, ops_per_sec}
|
||||
|
||||
pgxpool.New() — connection pool, MaxConns=20
|
||||
env: PGHOST, PGPORT, PGDATABASE, PGUSER, PGPASSWORD, PGSSLMODE
|
||||
```
|
||||
|
||||
### Цепочка таймаутов (после фиксов)
|
||||
|
||||
| Слой | Таймаут | Где задаётся |
|
||||
|------|---------|--------------|
|
||||
| curl --max-time | 720s | клиент |
|
||||
| nginx ingress proxy-read-timeout | 900s | аннотация ingress sless-operator |
|
||||
| operator http.Client (invoke.go) | TimeoutSec+5s = 705s | Function.Spec.TimeoutSec=700 |
|
||||
| function runtime (Go) | duration_sec = 600s | параметр в JSON body |
|
||||
|
||||
### Что было сломано до этой сессии
|
||||
|
||||
1. invoke.go: `var httpClient = &http.Client{Timeout: 30 * time.Second}` — хардкод в строке 24
|
||||
2. nginx ingress sless-operator: аннотации proxy-read-timeout не было → nginx дефолт 60s → 504
|
||||
|
||||
## 2026-03-19 — Tests 3-7: E2E прогон POSTGRES + 8 стресс-функций
|
||||
|
||||
### Tests 3-7
|
||||
|
||||
| Тест | Действие | Результат |
|
||||
|------|----------|-----------|
|
||||
| Test 3 | `table_rw.py`: добавлен `version: v2-with-hostname`, `host: socket.gethostname()` | ✅ Terraform apply, функция пересобрана |
|
||||
| Test 4 | `pg_info.js`: добавлен `code_version: v2-agent-test` | ✅ code_hash изменился, вернул `{"code_version":"v2-agent-test"}` |
|
||||
| Test 5 | Удалить trigger → 404 → пересоздать → работает | ✅ |
|
||||
| Test 6 | Удалить function + trigger → 404 → пересоздать (rebuild 46s) → работает | ✅ |
|
||||
| Test 7 | Новая функция `pg-stats` (Python): version, total_rows → протестирована → удалена | ✅ `{"version":"v1-test7","total_rows":3}` |
|
||||
|
||||
### 8 стресс-функций (коммит `014b99e`)
|
||||
|
||||
Задеплоены и прогнаны дважды через `stress_test.sh` (3 раунда).
|
||||
После crash-тестов все 8 подов: `Running`, 0 restarts.
|
||||
|
||||
| # | Функция | Runtime | Что проверяет | Результат |
|
||||
|---|---------|---------|---------------|-----------|
|
||||
| 1 | `stress-slow` | Python | sleep 3-N сек | ✅ `{"slept_sec":3}` |
|
||||
| 2 | `stress-bigloop` | Python | CPU n=2M | ✅ 0.31s |
|
||||
| 3 | `stress-divzero` | Python | ZeroDivisionError crash | ✅ pod restart → при d=7: `{"result":6.0}` |
|
||||
| 4 | `stress-writer` | Python | batch INSERT в PG | ✅ +3/+10 строк |
|
||||
| 5 | `stress-go-fast` | Go | factorial+fib без deps | ✅ factorial(20), fib(20) |
|
||||
| 6 | `stress-go-nil` | Go | nil pointer panic | ✅ pod restart → при crash=false: ok |
|
||||
| 7 | `stress-js-async` | NodeJS | 3 PG запроса Promise.all | ✅ version/count/max_id |
|
||||
| 8 | `stress-js-badenv` | NodeJS | TypeError (pod жив, HTTP 500) | ✅ нет restart |
|
||||
|
||||
**Итого в таблице после двух прогонов: 32 строки.**
|
||||
|
||||
### Поведение crash-функций
|
||||
|
||||
- `stress-divzero` / `stress-go-nil`: **pod restart** при краше (EOF на первый запрос после падения) — нормальное k8s поведение, runtime процесс умирает целиком
|
||||
- `stress-js-badenv`: **HTTP 500 без restart** — NodeJS поймал TypeError внутри async handler, pod остался жив
|
||||
- Это различие задокументировано: Python/Go crash = process exit, NodeJS crash = unhandled rejection в async = 500
|
||||
|
||||
### git
|
||||
|
||||
- Бинарник `event-dispatcher` (50MB) удалён из истории через `git reset --soft HEAD~2`
|
||||
- Добавлен в `.gitignore`
|
||||
- Force push: `d879817` → `014b99e`
|
||||
|
||||
---
|
||||
|
||||
## 2026-03-19 — event-trigger refactor (ветка feat/event-trigger-refactor)
|
||||
|
||||
| # | Задача | Статус | Заметки |
|
||||
|---|--------|--------|---------|
|
||||
| 1 | Аудит кластера — удаление мусора | ✅ | Удалены: event-monitor/writer/cleaner, pg-* тестовые функции, failed build jobs, orphan namespaces b1e4df2d/b794a3c4 |
|
||||
| 2 | Managed RabbitMQ через Terraform | ✅ | `terraform/RABBIT/`, ns `1dbfe9da-...`, host: `rabbitmqk8s.1dbfe9da-...svc.cluster.local` |
|
||||
| 3 | Архитектура event-trigger задокументирована | ✅ | `doc/decisions/log.md` — выбран Вариант A (отдельный event-dispatcher) |
|
||||
| 4 | trigger_types.go: +TriggerTypeEvent, +Queue | ✅ | |
|
||||
| 5 | internal/config: +RabbitMQURL | ✅ | |
|
||||
| 6 | services/event-dispatcher/ | ✅ | Go-сервис: dispatcher.go, watcher.go, main.go |
|
||||
| 7 | controllers/trigger_controller.go: +reconcileEvent | ✅ | Создаёт Service, обновляет status |
|
||||
| 8 | deployments/k8s/event-dispatcher.yaml | ✅ | ServiceAccount + ClusterRole + Deployment |
|
||||
| 9 | Образ naeel/sless-event-dispatcher:v0.1.0 | ✅ | Собран, запушен в DockerHub |
|
||||
| 10 | RABBITMQ_URL в sless-operator-secret | ✅ | managed rabbit (1dbfe9da namespace) |
|
||||
| 11 | event-dispatcher задеплоен в кластер | ✅ | Running в sless namespace |
|
||||
|
||||
|
||||
|
||||
## 2026-03-19 — pg-table-writer HTML + bugfix invoke.go Content-Length (оператор v0.1.37)
|
||||
|
||||
|
||||
@@ -36,6 +36,7 @@ exports.info = async (event) => {
|
||||
node_version: process.version,
|
||||
pg_version: versionRes.rows[0].v,
|
||||
table_rows: parseInt(countRes.rows[0].cnt, 10),
|
||||
code_version: 'v2-agent-test',
|
||||
};
|
||||
} finally {
|
||||
await client.end();
|
||||
|
||||
@@ -0,0 +1,38 @@
|
||||
# 2026-03-19
|
||||
# pg_stats.py — тестовая функция (Test 7): возвращает агрегированную статистику
|
||||
# по таблице terraform_demo_table: кол-во строк, дата первой и последней записи.
|
||||
# Создаётся и удаляется в рамках тестового прогона.
|
||||
#
|
||||
# Entrypoint: pg_stats.get_stats
|
||||
|
||||
import os
|
||||
import psycopg2
|
||||
import json
|
||||
|
||||
_CODE_VERSION = "v1-test7"
|
||||
|
||||
|
||||
def get_stats(event):
|
||||
conn = psycopg2.connect(
|
||||
host=os.environ["PGHOST"],
|
||||
port=int(os.environ.get("PGPORT", "5432")),
|
||||
dbname=os.environ["PGDATABASE"],
|
||||
user=os.environ["PGUSER"],
|
||||
password=os.environ["PGPASSWORD"],
|
||||
sslmode=os.environ.get("PGSSLMODE", "require"),
|
||||
)
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
cur.execute(
|
||||
"SELECT COUNT(*) AS cnt, MIN(created_at) AS first, MAX(created_at) AS last "
|
||||
"FROM terraform_demo_table"
|
||||
)
|
||||
row = cur.fetchone()
|
||||
return {
|
||||
"version": _CODE_VERSION,
|
||||
"total_rows": row[0],
|
||||
"first_row_at": str(row[1]) if row[1] else None,
|
||||
"last_row_at": str(row[2]) if row[2] else None,
|
||||
}
|
||||
finally:
|
||||
conn.close()
|
||||
@@ -0,0 +1 @@
|
||||
psycopg2-binary==2.9.9
|
||||
@@ -0,0 +1,20 @@
|
||||
# 2026-03-19
|
||||
# stress_bigloop.py — CPU-интенсивная функция: считает сумму квадратов N чисел.
|
||||
# Проверяет поведение под нагрузкой (большая и средняя итерация).
|
||||
|
||||
import time
|
||||
|
||||
_VERSION = "v1"
|
||||
|
||||
|
||||
def run(event):
|
||||
n = int(event.get("n", 500_000))
|
||||
start = time.monotonic()
|
||||
total = sum(i * i for i in range(n))
|
||||
elapsed = round(time.monotonic() - start, 4)
|
||||
return {
|
||||
"version": _VERSION,
|
||||
"n": n,
|
||||
"sum_of_squares": total,
|
||||
"elapsed_sec": elapsed,
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
# 2026-03-19
|
||||
# stress_divzero.py — намеренно делит на ноль (ZeroDivisionError).
|
||||
# Проверяет: платформа перехватывает панику, возвращает HTTP 500, не роняет под.
|
||||
|
||||
_VERSION = "v1"
|
||||
|
||||
|
||||
def run(event):
|
||||
numerator = int(event.get("n", 42))
|
||||
denominator = int(event.get("d", 0)) # по умолчанию 0 — намеренный краш
|
||||
# ZeroDivisionError: проверяем что платформа обрабатывает исключения
|
||||
result = numerator / denominator
|
||||
return {"version": _VERSION, "result": result}
|
||||
@@ -0,0 +1,43 @@
|
||||
package handler
|
||||
// 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),
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
// 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,
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,148 @@
|
||||
// 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),
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,7 @@
|
||||
{
|
||||
"name": "stress-js-async",
|
||||
"version": "1.0.0",
|
||||
"dependencies": {
|
||||
"pg": "^8.11.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,37 @@
|
||||
// 2026-03-19
|
||||
// stress_js_async.js — делает 3 параллельных запроса к PG через Promise.all.
|
||||
// Проверяет nodejs20 runtime под умеренной нагрузкой и async/await.
|
||||
//
|
||||
// Entrypoint: stress_js_async.run
|
||||
|
||||
'use strict';
|
||||
|
||||
const { Client } = require('pg');
|
||||
|
||||
exports.run = async (event) => {
|
||||
const client = new Client({
|
||||
host: process.env.PGHOST,
|
||||
port: parseInt(process.env.PGPORT || '5432'),
|
||||
database: process.env.PGDATABASE,
|
||||
user: process.env.PGUSER,
|
||||
password: process.env.PGPASSWORD,
|
||||
ssl: process.env.PGSSLMODE === 'require' ? { rejectUnauthorized: false } : false,
|
||||
});
|
||||
await client.connect();
|
||||
try {
|
||||
const [ver, cnt, max] = await Promise.all([
|
||||
client.query('SELECT version() AS v'),
|
||||
client.query('SELECT COUNT(*) AS cnt FROM terraform_demo_table'),
|
||||
client.query('SELECT MAX(id) AS max_id FROM terraform_demo_table'),
|
||||
]);
|
||||
return {
|
||||
runtime: 'nodejs20',
|
||||
version: 'v1',
|
||||
pg_version: ver.rows[0].v.split(' ').slice(0, 2).join(' '),
|
||||
total_rows: parseInt(cnt.rows[0].cnt, 10),
|
||||
max_id: max.rows[0].max_id,
|
||||
};
|
||||
} finally {
|
||||
await client.end();
|
||||
}
|
||||
};
|
||||
@@ -0,0 +1,5 @@
|
||||
{
|
||||
"name": "stress-js-badenv",
|
||||
"version": "1.0.0",
|
||||
"dependencies": {}
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
// 2026-03-19
|
||||
// stress_js_badenv.js — читает несуществующую переменную env и падает.
|
||||
// Проверяет: платформа перехватывает TypeError/undefined, возвращает 500.
|
||||
//
|
||||
// Entrypoint: stress_js_badenv.run
|
||||
|
||||
'use strict';
|
||||
|
||||
exports.run = async (event) => {
|
||||
const crash = event.crash !== false; // по умолчанию crash=true
|
||||
if (crash) {
|
||||
// Читаем несуществующий env, пытаемся вызвать .toUpperCase() на undefined
|
||||
const val = process.env.THIS_VAR_DOES_NOT_EXIST_AT_ALL;
|
||||
return { shout: val.toUpperCase() }; // TypeError: Cannot read properties of undefined
|
||||
}
|
||||
return { runtime: 'nodejs20', version: 'v1', crashed: false };
|
||||
};
|
||||
@@ -0,0 +1,18 @@
|
||||
# 2026-03-19
|
||||
# stress_slow.py — долгая функция: спит N секунд (по умолчанию 8).
|
||||
# Проверяет что timeout-механизм и параллельные запросы не блокируют друг друга.
|
||||
|
||||
import time
|
||||
import os
|
||||
|
||||
_VERSION = "v1"
|
||||
|
||||
|
||||
def run(event):
|
||||
secs = int(event.get("sleep", 8))
|
||||
time.sleep(secs)
|
||||
return {
|
||||
"version": _VERSION,
|
||||
"slept_sec": secs,
|
||||
"pid": os.getpid(),
|
||||
}
|
||||
@@ -0,0 +1 @@
|
||||
psycopg2-binary==2.9.9
|
||||
@@ -0,0 +1,39 @@
|
||||
# 2026-03-19
|
||||
# stress_writer.py — пишет N строк в terraform_demo_table (по умолчанию 5).
|
||||
# Проверяет параллельные INSERT'ы и устойчивость соединения с PG при нагрузке.
|
||||
|
||||
import os
|
||||
import psycopg2
|
||||
import time
|
||||
|
||||
_VERSION = "v1"
|
||||
|
||||
|
||||
def run(event):
|
||||
n = int(event.get("rows", 5))
|
||||
prefix = event.get("prefix", "stress")
|
||||
|
||||
conn = psycopg2.connect(
|
||||
host=os.environ["PGHOST"],
|
||||
port=int(os.environ.get("PGPORT", "5432")),
|
||||
dbname=os.environ["PGDATABASE"],
|
||||
user=os.environ["PGUSER"],
|
||||
password=os.environ["PGPASSWORD"],
|
||||
sslmode=os.environ.get("PGSSLMODE", "require"),
|
||||
)
|
||||
inserted = []
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
for i in range(n):
|
||||
title = f"{prefix}-{int(time.time()*1000)}-{i}"
|
||||
cur.execute(
|
||||
"INSERT INTO terraform_demo_table (title) VALUES (%s) RETURNING id",
|
||||
(title,),
|
||||
)
|
||||
row = cur.fetchone()
|
||||
inserted.append({"id": row[0], "title": title})
|
||||
conn.commit()
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
return {"version": _VERSION, "inserted": inserted, "count": len(inserted)}
|
||||
@@ -1,13 +1,16 @@
|
||||
# 2026-03-19
|
||||
# 2026-03-19 — добавлен version и hostname в ответ list_rows для тестирования обновления кода
|
||||
# table_rw.py — чтение и запись строк в terraform_demo_table.
|
||||
# Два entrypoint в одном файле: list_rows (JSON API) и add_row (HTML-страница + POST-обработчик).
|
||||
# ENV: PGHOST, PGPORT, PGDATABASE, PGUSER, PGPASSWORD, PGSSLMODE
|
||||
|
||||
import os
|
||||
import json
|
||||
import socket
|
||||
import psycopg2
|
||||
import psycopg2.extras
|
||||
|
||||
_CODE_VERSION = "v2-with-hostname"
|
||||
|
||||
|
||||
def _connect():
|
||||
return psycopg2.connect(
|
||||
@@ -29,7 +32,7 @@ def list_rows(event):
|
||||
"SELECT id, title, created_at::text FROM terraform_demo_table ORDER BY created_at DESC"
|
||||
)
|
||||
rows = [dict(r) for r in cur.fetchall()]
|
||||
return {"rows": rows, "count": len(rows)}
|
||||
return {"rows": rows, "count": len(rows), "version": _CODE_VERSION, "host": socket.gethostname()}
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
@@ -1,25 +1,23 @@
|
||||
|
||||
# resource "nubes_lucee" "app1" {
|
||||
# # Lucee-приложение, зависит от Postgres
|
||||
# # Lucee-приложение, подключается к nubes_postgres.npg (наш Postgres в реалме var.realm).
|
||||
# resource_name = "lucy_teststand_0"
|
||||
# # resource_realm = "k8s-3.ext.nubes.ru"
|
||||
# resource_realm = nubes_postgres.db2.resource_realm
|
||||
# # resource_realm = "k8s-4-sandbox-nubes-ru"
|
||||
# resource_realm = nubes_postgres.npg.resource_realm
|
||||
# domain = "web-test-stand"
|
||||
|
||||
# git_path = "https://gitea-naeel.giteak8s.services.ngcloud.ru/naeel/testlucee"
|
||||
# git_path = "https://gitea-naeel.giteak8s.services.ngcloud.ru/naeel/testlucee"
|
||||
|
||||
# json_env = jsonencode({
|
||||
# # 🔗 Настройки Data Source 'testds' для Lucee (Application.cfc)
|
||||
# testds_class = "org.postgresql.Driver" # 📂 Драйвер БД
|
||||
# testds_bundleName = "org.postgresql.jdbc" # 📦 Имя бандла JDBC
|
||||
# testds_bundleVersion = "42.6.0" # 🔢 Версия драйвера
|
||||
# testds_connectionString = "jdbc:postgresql://${nubes_postgres.db2.state_out_flat["internalConnect.master"]}:5432/postgres?sslmode=require" # 🚀 Строка подключения
|
||||
# testds_username = nubes_postgres_user.db2_user.username # 👤 Логин
|
||||
# testds_password = jsondecode(nubes_postgres.db2.vault_secrets["users"])[nubes_postgres_user.db2_user.username]["password"] # 🔑 Пароль
|
||||
# testds_connectionLimit = "5" # 🚦 Лимит соединений
|
||||
# testds_liveTimeout = "15" # ⏳ Таймаут жизни
|
||||
# testds_validate = "false" # ✅ Валидация при запросе
|
||||
# json_env = jsonencode({
|
||||
# # Настройки Data Source 'testds' для Lucee (Application.cfc)
|
||||
# testds_class = "org.postgresql.Driver"
|
||||
# testds_bundleName = "org.postgresql.jdbc"
|
||||
# testds_bundleVersion = "42.6.0"
|
||||
# testds_connectionString = "jdbc:postgresql://${local.pg_host}:5432/${local.pg_database}?sslmode=require"
|
||||
# testds_username = local.pg_username
|
||||
# testds_password = local.pg_password
|
||||
# testds_connectionLimit = "5"
|
||||
# testds_liveTimeout = "15"
|
||||
# testds_validate = "false"
|
||||
# })
|
||||
|
||||
# resource_c_p_u = 300
|
||||
@@ -27,40 +25,40 @@
|
||||
# resource_instances = 1
|
||||
# app_version = "5.4"
|
||||
|
||||
# depends_on = [nubes_postgres.db2]
|
||||
# depends_on = [nubes_postgres.npg]
|
||||
# }
|
||||
|
||||
# resource "nubes_nodejs" "app3" {
|
||||
# # NodeJS демо, работающий с тем же Postgres.
|
||||
# resource_name = "node_01"
|
||||
# resource_realm = nubes_postgres.db2.resource_realm
|
||||
# domain = "node07"
|
||||
# git_path = "https://gitea-naeel.giteak8s.services.ngcloud.ru/naeel/testnode.git"
|
||||
# health_path = "/healthz"
|
||||
# app_version = "23"
|
||||
|
||||
# json_env = jsonencode({
|
||||
# # Переменные подключения к Postgres.
|
||||
# PGHOST = nubes_postgres.db2.state_out_flat["internalConnect.master"]
|
||||
# PGPORT = "5432"
|
||||
# PGUSER = nubes_postgres_user.db2_user.username
|
||||
# PGPASSWORD = jsondecode(nubes_postgres.db2.vault_secrets["users"])[nubes_postgres_user.db2_user.username]["password"]
|
||||
# PGDATABASE = nubes_postgres_database.db2_app.db_name
|
||||
# PGSSLMODE = "require"
|
||||
# DATABASE_URL = format(
|
||||
# # NodeJS — застрял при создании на тест-стенде (операция в ожидании).
|
||||
# # Раскомментировать когда тест-стенд стабилен.
|
||||
# resource_name = "node_01"
|
||||
# resource_realm = nubes_postgres.npg.resource_realm
|
||||
# domain = "node07"
|
||||
# git_path = "https://gitea-naeel.giteak8s.services.ngcloud.ru/naeel/testnode.git"
|
||||
# health_path = "/healthz"
|
||||
# app_version = "23"
|
||||
#
|
||||
# json_env = jsonencode({
|
||||
# PGHOST = local.pg_host
|
||||
# PGPORT = "5432"
|
||||
# PGUSER = local.pg_username
|
||||
# PGPASSWORD = local.pg_password
|
||||
# PGDATABASE = local.pg_database
|
||||
# PGSSLMODE = "require"
|
||||
# DATABASE_URL = format(
|
||||
# "postgresql://%s:%s@%s:5432/%s?sslmode=require",
|
||||
# nubes_postgres_user.db2_user.username,
|
||||
# jsondecode(nubes_postgres.db2.vault_secrets["users"])[nubes_postgres_user.db2_user.username]["password"],
|
||||
# nubes_postgres.db2.state_out_flat["internalConnect.master"],
|
||||
# nubes_postgres_database.db2_app.db_name
|
||||
# )
|
||||
# })
|
||||
|
||||
# resource_c_p_u = 300
|
||||
# resource_memory = 256
|
||||
# resource_instances = 1
|
||||
|
||||
# depends_on = [nubes_postgres.db2]
|
||||
# local.pg_username,
|
||||
# local.pg_password,
|
||||
# local.pg_host,
|
||||
# local.pg_database
|
||||
# )
|
||||
# })
|
||||
#
|
||||
# resource_c_p_u = 300
|
||||
# resource_memory = 256
|
||||
# resource_instances = 1
|
||||
#
|
||||
# depends_on = [nubes_postgres.npg]
|
||||
# }
|
||||
|
||||
# output "pg_vault_secrets" {
|
||||
@@ -68,6 +66,4 @@
|
||||
# sensitive = true
|
||||
# }
|
||||
|
||||
# terraform output -json pg_vault_secrets
|
||||
|
||||
|
||||
|
||||
@@ -47,12 +47,12 @@ variable "pg_password" {
|
||||
|
||||
provider "nubes" {
|
||||
api_token = var.api_token
|
||||
api_endpoint = "https://deck-api-test.ngcloud.ru/api/v1/index.cfm"
|
||||
api_endpoint = "https://deck-test.ngcloud.ru/api/v1/index.cfm"
|
||||
}
|
||||
|
||||
provider "sless" {
|
||||
endpoint = "https://sless.kube5s.ru"
|
||||
token = var.api_token
|
||||
nubes_endpoint = "https://deck-api-test.ngcloud.ru/api/v1"
|
||||
nubes_endpoint = "https://deck-test.ngcloud.ru/api/v1"
|
||||
}
|
||||
|
||||
|
||||
@@ -166,8 +166,8 @@ resource "sless_function" "postgres_table_writer" {
|
||||
name = "pg-table-writer"
|
||||
runtime = "python3.11"
|
||||
entrypoint = "table_rw.add_row"
|
||||
memory_mb = 128
|
||||
timeout_sec = 30
|
||||
memory_mb = 256
|
||||
timeout_sec = 45
|
||||
|
||||
env_vars = {
|
||||
PGHOST = local.pg_host
|
||||
@@ -194,3 +194,214 @@ output "table_writer_url" {
|
||||
value = sless_trigger.postgres_table_writer_http.url
|
||||
}
|
||||
|
||||
# =============================================================================
|
||||
# STRESS-ТЕСТЫ: 8 функций для проверки устойчивости платформы.
|
||||
# Python: slow, divzero, bigloop, writer
|
||||
# Go: fast, nil-panic
|
||||
# NodeJS: async-parallel, badenv
|
||||
# =============================================================================
|
||||
|
||||
# --- [1] Python: долгая (sleep N сек) ---
|
||||
resource "sless_function" "stress_slow" {
|
||||
name = "stress-slow"
|
||||
runtime = "python3.11"
|
||||
entrypoint = "stress_slow.run"
|
||||
memory_mb = 128
|
||||
timeout_sec = 30
|
||||
|
||||
env_vars = {}
|
||||
|
||||
source_dir = "${path.module}/code/stress-slow"
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
resource "sless_trigger" "stress_slow_http" {
|
||||
name = "stress-slow-http"
|
||||
type = "http"
|
||||
function = sless_function.stress_slow.name
|
||||
enabled = true
|
||||
}
|
||||
|
||||
# --- [2] Python: деление на ноль ---
|
||||
resource "sless_function" "stress_divzero" {
|
||||
name = "stress-divzero"
|
||||
runtime = "python3.11"
|
||||
entrypoint = "stress_divzero.run"
|
||||
memory_mb = 128
|
||||
timeout_sec = 10
|
||||
|
||||
env_vars = {}
|
||||
|
||||
source_dir = "${path.module}/code/stress-divzero"
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
resource "sless_trigger" "stress_divzero_http" {
|
||||
name = "stress-divzero-http"
|
||||
type = "http"
|
||||
function = sless_function.stress_divzero.name
|
||||
enabled = true
|
||||
}
|
||||
|
||||
# --- [3] Python: CPU bigloop ---
|
||||
resource "sless_function" "stress_bigloop" {
|
||||
name = "stress-bigloop"
|
||||
runtime = "python3.11"
|
||||
entrypoint = "stress_bigloop.run"
|
||||
memory_mb = 256
|
||||
timeout_sec = 30
|
||||
|
||||
env_vars = {}
|
||||
|
||||
source_dir = "${path.module}/code/stress-bigloop"
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
resource "sless_trigger" "stress_bigloop_http" {
|
||||
name = "stress-bigloop-http"
|
||||
type = "http"
|
||||
function = sless_function.stress_bigloop.name
|
||||
enabled = true
|
||||
}
|
||||
|
||||
# --- [4] Python: массовая запись в PG ---
|
||||
resource "sless_function" "stress_writer" {
|
||||
name = "stress-writer"
|
||||
runtime = "python3.11"
|
||||
entrypoint = "stress_writer.run"
|
||||
memory_mb = 128
|
||||
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/stress-writer"
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
resource "sless_trigger" "stress_writer_http" {
|
||||
name = "stress-writer-http"
|
||||
type = "http"
|
||||
function = sless_function.stress_writer.name
|
||||
enabled = true
|
||||
}
|
||||
|
||||
# --- [5] Go: быстрая математика (факториал + Фибоначчи) ---
|
||||
resource "sless_function" "stress_go_fast" {
|
||||
name = "stress-go-fast"
|
||||
runtime = "go1.23"
|
||||
entrypoint = "handler.Handle"
|
||||
memory_mb = 64
|
||||
timeout_sec = 10
|
||||
|
||||
env_vars = {}
|
||||
|
||||
source_dir = "${path.module}/code/stress-go-fast"
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
resource "sless_trigger" "stress_go_fast_http" {
|
||||
name = "stress-go-fast-http"
|
||||
type = "http"
|
||||
function = sless_function.stress_go_fast.name
|
||||
enabled = true
|
||||
}
|
||||
|
||||
# --- [6] Go: nil pointer panic ---
|
||||
resource "sless_function" "stress_go_nil" {
|
||||
name = "stress-go-nil"
|
||||
runtime = "go1.23"
|
||||
entrypoint = "handler.Handle"
|
||||
memory_mb = 64
|
||||
timeout_sec = 10
|
||||
|
||||
env_vars = {}
|
||||
|
||||
source_dir = "${path.module}/code/stress-go-nil"
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
resource "sless_trigger" "stress_go_nil_http" {
|
||||
name = "stress-go-nil-http"
|
||||
type = "http"
|
||||
function = sless_function.stress_go_nil.name
|
||||
enabled = true
|
||||
}
|
||||
|
||||
# --- [7] NodeJS: 3 параллельных запроса к PG через Promise.all ---
|
||||
resource "sless_function" "stress_js_async" {
|
||||
name = "stress-js-async"
|
||||
runtime = "nodejs20"
|
||||
entrypoint = "stress_js_async.run"
|
||||
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/stress-js-async"
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
resource "sless_trigger" "stress_js_async_http" {
|
||||
name = "stress-js-async-http"
|
||||
type = "http"
|
||||
function = sless_function.stress_js_async.name
|
||||
enabled = true
|
||||
}
|
||||
|
||||
# --- [8] NodeJS: TypeError на несуществующей env-переменной ---
|
||||
resource "sless_function" "stress_js_badenv" {
|
||||
name = "stress-js-badenv"
|
||||
runtime = "nodejs20"
|
||||
entrypoint = "stress_js_badenv.run"
|
||||
memory_mb = 128
|
||||
timeout_sec = 10
|
||||
|
||||
env_vars = {}
|
||||
|
||||
source_dir = "${path.module}/code/stress-js-badenv"
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
resource "sless_trigger" "stress_js_badenv_http" {
|
||||
name = "stress-js-badenv-http"
|
||||
type = "http"
|
||||
function = sless_function.stress_js_badenv.name
|
||||
enabled = true
|
||||
}
|
||||
|
||||
# --- [9] Go: PG Storm — 100 горутин долбят PostgreSQL напрямую через pgxpool ---
|
||||
# Тестирует: Go runtime под конкурентной нагрузкой, pgxpool connection pool,
|
||||
# устойчивость managed PG при массовых INSERT/SELECT.
|
||||
# Параметры: workers (default 100), duration_sec (default 600), max_delay_ms (default 300).
|
||||
resource "sless_function" "stress_go_pgstorm" {
|
||||
name = "stress-go-pgstorm"
|
||||
runtime = "go1.23"
|
||||
entrypoint = "handler.Handle"
|
||||
memory_mb = 256
|
||||
timeout_sec = 700
|
||||
|
||||
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/stress-go-pgstorm"
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
resource "sless_trigger" "stress_go_pgstorm_http" {
|
||||
name = "stress-go-pgstorm-http"
|
||||
type = "http"
|
||||
function = sless_function.stress_go_pgstorm.name
|
||||
enabled = true
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,57 @@
|
||||
#!/bin/bash
|
||||
# 2026-03-19 — stress test script: параллельный запуск всех 8 стресс-функций
|
||||
BASE="https://sless.kube5s.ru/fn/sless-ffd1f598c169b0ae"
|
||||
|
||||
echo "=== РАУНД 1: первый холодный запуск ==="
|
||||
curl -s -m 35 "$BASE/stress-slow" -d '{"sleep":3}' -H "Content-Type:application/json" > /tmp/r_slow.json &
|
||||
curl -s -m 10 "$BASE/stress-divzero" > /tmp/r_divzero.json &
|
||||
curl -s -m 40 "$BASE/stress-bigloop" -d '{"n":1000000}' -H "Content-Type:application/json"> /tmp/r_bigloop.json &
|
||||
curl -s -m 35 "$BASE/stress-writer" -d '{"rows":3,"prefix":"batch1"}' -H "Content-Type:application/json" > /tmp/r_writer.json &
|
||||
curl -s -m 15 "$BASE/stress-go-fast" -d '{"n":15}' -H "Content-Type:application/json" > /tmp/r_go_fast.json &
|
||||
curl -s -m 10 "$BASE/stress-go-nil" > /tmp/r_go_nil.json &
|
||||
curl -s -m 20 "$BASE/stress-js-async" > /tmp/r_js_async.json &
|
||||
curl -s -m 10 "$BASE/stress-js-badenv" > /tmp/r_js_badenv.json &
|
||||
wait
|
||||
|
||||
echo "[slow]: $(cat /tmp/r_slow.json)"
|
||||
echo "[divzero]: $(cat /tmp/r_divzero.json)"
|
||||
echo "[bigloop]: $(cat /tmp/r_bigloop.json)"
|
||||
echo "[writer]: $(cat /tmp/r_writer.json)"
|
||||
echo "[go-fast]: $(cat /tmp/r_go_fast.json)"
|
||||
echo "[go-nil]: $(cat /tmp/r_go_nil.json)"
|
||||
echo "[js-async]: $(cat /tmp/r_js_async.json)"
|
||||
echo "[js-badenv]:$(cat /tmp/r_js_badenv.json)"
|
||||
|
||||
echo ""
|
||||
echo "=== РАУНД 2: повторный (горячий кэш) ==="
|
||||
curl -s -m 15 "$BASE/stress-bigloop" -d '{"n":2000000}' -H "Content-Type:application/json" > /tmp/r2_bigloop.json &
|
||||
curl -s -m 10 "$BASE/stress-go-fast" -d '{"n":20}' -H "Content-Type:application/json" > /tmp/r2_go_fast.json &
|
||||
curl -s -m 20 "$BASE/stress-js-async" > /tmp/r2_async.json &
|
||||
curl -s -m 35 "$BASE/stress-writer" -d '{"rows":10,"prefix":"batch2"}' -H "Content-Type:application/json" > /tmp/r2_writer.json &
|
||||
wait
|
||||
echo "[bigloop-2M]: $(cat /tmp/r2_bigloop.json)"
|
||||
echo "[go-fast-20]: $(cat /tmp/r2_go_fast.json)"
|
||||
echo "[js-async-2]: $(cat /tmp/r2_async.json)"
|
||||
echo "[writer-10]: $(cat /tmp/r2_writer.json)"
|
||||
|
||||
echo ""
|
||||
echo "=== РАУНД 3: crash функции с неверными параметрами ==="
|
||||
curl -s -m 10 "$BASE/stress-divzero" -d '{"n":100,"d":0}' -H "Content-Type:application/json" > /tmp/r3_dz.json &
|
||||
curl -s -m 10 "$BASE/stress-go-nil" -d '{"crash":true}' -H "Content-Type:application/json" > /tmp/r3_nil.json &
|
||||
curl -s -m 10 "$BASE/stress-js-badenv" -d '{"crash":true}' -H "Content-Type:application/json" > /tmp/r3_bad.json &
|
||||
# divzero с нормальным делителем — должен вернуть результат
|
||||
curl -s -m 10 "$BASE/stress-divzero" -d '{"n":42,"d":7}' -H "Content-Type:application/json" > /tmp/r3_ok.json &
|
||||
# go-nil без краша — должен вернуть ok
|
||||
curl -s -m 10 "$BASE/stress-go-nil" -d '{"crash":false}' -H "Content-Type:application/json" > /tmp/r3_nil_ok.json &
|
||||
wait
|
||||
echo "[divzero crash]: $(cat /tmp/r3_dz.json)"
|
||||
echo "[go-nil crash]: $(cat /tmp/r3_nil.json)"
|
||||
echo "[js-badenv crash]: $(cat /tmp/r3_bad.json)"
|
||||
echo "[divzero ok 42/7]: $(cat /tmp/r3_ok.json)"
|
||||
echo "[go-nil ok]: $(cat /tmp/r3_nil_ok.json)"
|
||||
|
||||
echo ""
|
||||
echo "=== ИТОГ: количество строк в таблице ==="
|
||||
curl -s -m 15 "$BASE/pg-table-reader"
|
||||
echo ""
|
||||
echo "=== DONE ==="
|
||||
@@ -0,0 +1,94 @@
|
||||
# 2026-03-18 (обновлено: plain text вывод; фильтрация SLESS_EXCLUDE)
|
||||
# funcs_list.py — HTTP-функция: список пользовательских функций, человекочитаемый plain text.
|
||||
# Вызывает внутренний REST API оператора (ClusterIP, без TLS).
|
||||
# Возвращает str → python runtime отдаёт text/plain напрямую без json.dumps.
|
||||
#
|
||||
# Env vars:
|
||||
# SLESS_API_URL — URL оператора (http://sless-operator.sless.svc.cluster.local:9090)
|
||||
# SLESS_NAMESPACE — namespace пользователя (sless-{hex16})
|
||||
# SLESS_TOKEN — JWT токен для /v1/ API
|
||||
# SLESS_EXTERNAL_URL — публичный базовый URL (https://sless.kube5s.ru)
|
||||
# SLESS_EXCLUDE — comma-separated имена функций, которые не показывать
|
||||
|
||||
import os
|
||||
import requests
|
||||
|
||||
SEP = "─" * 52
|
||||
|
||||
|
||||
def _comment(fn, http_trigs, cron_trigs):
|
||||
phase = fn.get("phase", "?")
|
||||
runtime = fn.get("runtime", "?")
|
||||
if http_trigs:
|
||||
active = "активна" if http_trigs[0].get("active") else "неактивна"
|
||||
return f"HTTP endpoint ({runtime}) — {phase}, {active}"
|
||||
elif cron_trigs:
|
||||
schedule = cron_trigs[0].get("schedule", "?")
|
||||
active = "активна" if cron_trigs[0].get("active") else "неактивна"
|
||||
return f"Cron '{schedule}' ({runtime}) — {phase}, {active}"
|
||||
else:
|
||||
return f"Job/runner без триггера ({runtime}) — {phase}"
|
||||
|
||||
|
||||
def list_all(event):
|
||||
api_url = os.environ["SLESS_API_URL"].rstrip("/")
|
||||
namespace = os.environ["SLESS_NAMESPACE"]
|
||||
token = os.environ["SLESS_TOKEN"]
|
||||
ext_url = os.environ.get("SLESS_EXTERNAL_URL", "").rstrip("/")
|
||||
exclude = {n.strip() for n in os.environ.get("SLESS_EXCLUDE", "").split(",") if n.strip()}
|
||||
|
||||
headers = {"Authorization": f"Bearer {token}"}
|
||||
fns = requests.get(f"{api_url}/v1/namespaces/{namespace}/functions", headers=headers, timeout=10)
|
||||
trs = requests.get(f"{api_url}/v1/namespaces/{namespace}/triggers", headers=headers, timeout=10)
|
||||
fns.raise_for_status()
|
||||
trs.raise_for_status()
|
||||
|
||||
trig_idx = {}
|
||||
for tr in trs.json():
|
||||
fn_name = tr.get("function") or tr.get("functionRef")
|
||||
if fn_name:
|
||||
trig_idx.setdefault(fn_name, []).append(tr)
|
||||
|
||||
items = []
|
||||
for fn in fns.json():
|
||||
name = fn["name"]
|
||||
if name in exclude:
|
||||
continue
|
||||
http_t = [t for t in trig_idx.get(name, []) if t.get("type") == "http"]
|
||||
cron_t = [t for t in trig_idx.get(name, []) if t.get("type") == "cron"]
|
||||
is_active = any(t.get("enabled", True) and t.get("active", False) for t in trig_idx.get(name, []))
|
||||
items.append((fn, http_t, cron_t, is_active))
|
||||
|
||||
# Сортировка: активные вверх, затем по имени
|
||||
items.sort(key=lambda x: (not x[3], x[0]["name"]))
|
||||
|
||||
lines = []
|
||||
for fn, http_t, cron_t, is_active in items:
|
||||
name = fn["name"]
|
||||
lines.append(SEP)
|
||||
lines.append(f" {_comment(fn, http_t, cron_t)}")
|
||||
lines.append(f" name: {name}")
|
||||
lines.append(f" runtime: {fn.get('runtime', '?')}")
|
||||
lines.append(f" phase: {fn.get('phase', '?')}")
|
||||
lines.append(f" active: {'да' if is_active else 'нет'}")
|
||||
|
||||
if http_t:
|
||||
url = f"{ext_url}/fn/{namespace}/{name}" if ext_url else http_t[0].get("url", "")
|
||||
lines.append(f" url: {url}")
|
||||
if cron_t:
|
||||
lines.append(f" cron: {cron_t[0].get('schedule', '?')}")
|
||||
if fn.get("created_at"):
|
||||
lines.append(f" created: {fn['created_at']}")
|
||||
if fn.get("last_built_at"):
|
||||
lines.append(f" built: {fn['last_built_at']}")
|
||||
if fn.get("message"):
|
||||
lines.append(f" message: {fn['message']}")
|
||||
|
||||
lines.append(SEP)
|
||||
lines.append(f" namespace: {namespace} | total: {len(items)}")
|
||||
lines.append(SEP)
|
||||
|
||||
# Возвращаем str — python runtime отдаст text/plain напрямую
|
||||
return "\n".join(lines) + "\n"
|
||||
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
requests==2.31.0
|
||||
@@ -0,0 +1,8 @@
|
||||
{
|
||||
"name": "pg-info",
|
||||
"version": "1.0.0",
|
||||
"description": "sless nodejs20 function: pg version + table info",
|
||||
"dependencies": {
|
||||
"pg": "8.11.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,43 @@
|
||||
// 2026-03-18
|
||||
// pg_info.js — NodeJS-функция: проверка работы JS runtime + чтение мета-данных БД.
|
||||
// Подключается к PostgreSQL через пакет pg, возвращает версию сервера и счётчик строк.
|
||||
// Демонстрирует: nodejs20 runtime, npm-зависимость (package.json), PG из JS.
|
||||
//
|
||||
// ENV (те же что у python-функций):
|
||||
// PGHOST, PGPORT, PGDATABASE, PGUSER, PGPASSWORD, PGSSLMODE
|
||||
//
|
||||
// Entrypoint: pg_info.info
|
||||
|
||||
'use strict';
|
||||
|
||||
const { Client } = require('pg');
|
||||
|
||||
exports.info = async (event) => {
|
||||
const client = new Client({
|
||||
host: process.env.PGHOST,
|
||||
port: parseInt(process.env.PGPORT || '5432'),
|
||||
database: process.env.PGDATABASE,
|
||||
user: process.env.PGUSER,
|
||||
password: process.env.PGPASSWORD,
|
||||
// pg-пакет требует явного ssl-объекта; rejectUnauthorized: false — т.к.
|
||||
// self-signed cert на nubes managed PG, но канал всё равно шифруется.
|
||||
ssl: process.env.PGSSLMODE === 'require' ? { rejectUnauthorized: false } : false,
|
||||
});
|
||||
|
||||
await client.connect();
|
||||
try {
|
||||
const [versionRes, countRes] = await Promise.all([
|
||||
client.query('SELECT version() AS v'),
|
||||
client.query('SELECT COUNT(*) AS cnt FROM terraform_demo_table'),
|
||||
]);
|
||||
|
||||
return {
|
||||
runtime: 'nodejs20',
|
||||
node_version: process.version,
|
||||
pg_version: versionRes.rows[0].v,
|
||||
table_rows: parseInt(countRes.rows[0].cnt, 10),
|
||||
};
|
||||
} finally {
|
||||
await client.end();
|
||||
}
|
||||
};
|
||||
@@ -0,0 +1,3 @@
|
||||
# 2026-03-17 00:00
|
||||
# requirements.txt — зависимости для функции запуска SQL.
|
||||
psycopg2-binary==2.9.9
|
||||
@@ -0,0 +1,39 @@
|
||||
# 2026-03-17 00:00
|
||||
# sql_runner.py — функция для выполнения SQL-операторов из входного события.
|
||||
import os
|
||||
import psycopg2
|
||||
|
||||
|
||||
def run_sql(event):
|
||||
# Выполняет список SQL-операторов в одной транзакции для атомарной инициализации схемы.
|
||||
# Параметры подключения передаются раздельно, чтобы избежать ошибок парсинга DSN при спецсимволах.
|
||||
pg_host = os.environ["PGHOST"]
|
||||
pg_port = os.environ.get("PGPORT", "5432")
|
||||
pg_database = os.environ["PGDATABASE"]
|
||||
pg_user = os.environ["PGUSER"]
|
||||
pg_password = os.environ["PGPASSWORD"]
|
||||
pg_sslmode = os.environ.get("PGSSLMODE", "require")
|
||||
statements = event.get("statements", [])
|
||||
|
||||
if not statements:
|
||||
return {"error": "no statements provided"}
|
||||
|
||||
connection = psycopg2.connect(
|
||||
host=pg_host,
|
||||
port=pg_port,
|
||||
dbname=pg_database,
|
||||
user=pg_user,
|
||||
password=pg_password,
|
||||
sslmode=pg_sslmode,
|
||||
)
|
||||
try:
|
||||
cursor = connection.cursor()
|
||||
for statement in statements:
|
||||
cursor.execute(statement)
|
||||
connection.commit()
|
||||
return {"ok": True, "executed": len(statements)}
|
||||
except Exception as error:
|
||||
connection.rollback()
|
||||
return {"error": str(error)}
|
||||
finally:
|
||||
connection.close()
|
||||
@@ -0,0 +1 @@
|
||||
psycopg2-binary==2.9.9
|
||||
@@ -0,0 +1,130 @@
|
||||
# 2026-03-19
|
||||
# table_rw.py — чтение и запись строк в terraform_demo_table.
|
||||
# Два entrypoint в одном файле: list_rows (JSON API) и add_row (HTML-страница + POST-обработчик).
|
||||
# ENV: PGHOST, PGPORT, PGDATABASE, PGUSER, PGPASSWORD, PGSSLMODE
|
||||
|
||||
import os
|
||||
import json
|
||||
import psycopg2
|
||||
import psycopg2.extras
|
||||
|
||||
|
||||
def _connect():
|
||||
return psycopg2.connect(
|
||||
host=os.environ["PGHOST"],
|
||||
port=os.environ.get("PGPORT", "5432"),
|
||||
dbname=os.environ["PGDATABASE"],
|
||||
user=os.environ["PGUSER"],
|
||||
password=os.environ["PGPASSWORD"],
|
||||
sslmode=os.environ.get("PGSSLMODE", "require"),
|
||||
)
|
||||
|
||||
|
||||
def list_rows(event):
|
||||
# Возвращает все строки terraform_demo_table, отсортированные по убыванию created_at.
|
||||
conn = _connect()
|
||||
try:
|
||||
cur = conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor)
|
||||
cur.execute(
|
||||
"SELECT id, title, created_at::text FROM terraform_demo_table ORDER BY created_at DESC"
|
||||
)
|
||||
rows = [dict(r) for r in cur.fetchall()]
|
||||
return {"rows": rows, "count": len(rows)}
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def _render_page(rows, message=""):
|
||||
# HTML-страница с формой ввода и таблицей строк.
|
||||
# message — статус последней операции (успех / ошибка).
|
||||
rows_html = "".join(
|
||||
f"<tr><td>{r['id']}</td><td>{r['title']}</td><td>{r['created_at']}</td></tr>"
|
||||
for r in rows
|
||||
)
|
||||
msg_html = f'<p class="msg">{message}</p>' if message else ""
|
||||
return f"""<!DOCTYPE html>
|
||||
<html lang="ru">
|
||||
<head>
|
||||
<meta charset="utf-8">
|
||||
<title>pg-table-writer</title>
|
||||
<style>
|
||||
body {{ font-family: sans-serif; max-width: 700px; margin: 40px auto; background: #111; color: #eee; }}
|
||||
h1 {{ color: #7dd3fc; }}
|
||||
form {{ display: flex; gap: 8px; margin-bottom: 24px; }}
|
||||
input[type=text] {{ flex: 1; padding: 8px 12px; border-radius: 6px; border: 1px solid #444; background: #1e1e1e; color: #eee; font-size: 15px; }}
|
||||
button {{ padding: 8px 18px; background: #2563eb; color: #fff; border: none; border-radius: 6px; cursor: pointer; font-size: 15px; }}
|
||||
button:hover {{ background: #1d4ed8; }}
|
||||
table {{ width: 100%; border-collapse: collapse; }}
|
||||
th, td {{ padding: 8px 10px; border-bottom: 1px solid #333; text-align: left; }}
|
||||
th {{ color: #7dd3fc; }}
|
||||
.msg {{ color: #4ade80; margin-bottom: 12px; }}
|
||||
</style>
|
||||
</head>
|
||||
<body>
|
||||
<h1>pg-table-writer</h1>
|
||||
<form method="POST">
|
||||
<input type="text" name="title" placeholder="Введите строку..." autofocus required>
|
||||
<button type="submit">Добавить</button>
|
||||
</form>
|
||||
{msg_html}
|
||||
<table>
|
||||
<thead><tr><th>#</th><th>title</th><th>created_at</th></tr></thead>
|
||||
<tbody>{rows_html}</tbody>
|
||||
</table>
|
||||
</body>
|
||||
</html>"""
|
||||
|
||||
|
||||
def add_row(event):
|
||||
# GET → HTML-страница с формой и списком строк.
|
||||
# POST → вставляет строку из form-поля title или JSON-поля title,
|
||||
# затем возвращает обновлённую HTML-страницу.
|
||||
# POST с Content-Type: application/json (curl/API) → возвращает JSON.
|
||||
method = event.get("_method", "GET")
|
||||
|
||||
if method == "GET":
|
||||
conn = _connect()
|
||||
try:
|
||||
cur = conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor)
|
||||
cur.execute("SELECT id, title, created_at::text FROM terraform_demo_table ORDER BY created_at DESC")
|
||||
rows = [dict(r) for r in cur.fetchall()]
|
||||
finally:
|
||||
conn.close()
|
||||
return _render_page(rows)
|
||||
|
||||
# POST — вставка строки
|
||||
# Поле title приходит либо из JSON-тела, либо из application/x-www-form-urlencoded.
|
||||
# Сервер уже распарсил JSON в event; form-данные приходят как event["body"] = "title=...".
|
||||
title = event.get("title", "").strip()
|
||||
if not title:
|
||||
# Попытка распарсить form-encoded body (браузерная форма)
|
||||
body = event.get("body", "")
|
||||
if body.startswith("title="):
|
||||
from urllib.parse import unquote_plus
|
||||
title = unquote_plus(body[len("title="):].split("&")[0]).strip()
|
||||
|
||||
if not title:
|
||||
return {"ok": False, "error": "title is required"}
|
||||
|
||||
conn = _connect()
|
||||
try:
|
||||
cur = conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor)
|
||||
cur.execute(
|
||||
"INSERT INTO terraform_demo_table (title) VALUES (%s) RETURNING id, title, created_at::text",
|
||||
(title,),
|
||||
)
|
||||
row = dict(cur.fetchone())
|
||||
conn.commit()
|
||||
|
||||
# Если запрос из браузера (form POST) — возвращаем обновлённую страницу.
|
||||
# Если из curl/API — возвращаем JSON.
|
||||
accept = event.get("_accept", "")
|
||||
if "application/json" in accept:
|
||||
return {"ok": True, "row": row}
|
||||
|
||||
# Перечитываем все строки для обновлённой страницы
|
||||
cur.execute("SELECT id, title, created_at::text FROM terraform_demo_table ORDER BY created_at DESC")
|
||||
rows = [dict(r) for r in cur.fetchall()]
|
||||
return _render_page(rows, message=f"Добавлено: «{row['title']}»")
|
||||
finally:
|
||||
conn.close()
|
||||
@@ -0,0 +1,129 @@
|
||||
# 2026-03-18 (обновлено: фильтрация SLESS_EXCLUDE, читаемый вывод через "#"-ключ)
|
||||
# funcs_list.py — HTTP-функция: список всех пользовательских функций с их статусами.
|
||||
# Вызывает внутренний REST API оператора (ClusterIP, без TLS).
|
||||
# Объединяет данные функций и триггеров в один ответ; скрывает служебные функции.
|
||||
#
|
||||
# Env vars:
|
||||
# SLESS_API_URL — URL оператора (http://sless-operator.sless.svc.cluster.local:9090)
|
||||
# SLESS_NAMESPACE — namespace пользователя (sless-{hex16})
|
||||
# SLESS_TOKEN — JWT токен для /v1/ API
|
||||
# SLESS_EXTERNAL_URL — публичный базовый URL (https://sless.kube5s.ru), для корректных ссылок
|
||||
# SLESS_EXCLUDE — comma-separated имена функций, которые не надо показывать
|
||||
# Пример: "funcs,event-writer,event-monitor,event-cleaner"
|
||||
#
|
||||
# Формат вывода: JSON-объект, где каждая функция содержит поле "#" — краткий комментарий.
|
||||
# При pretty-print (python3 -m json.tool) выглядит как читаемый список с аннотациями.
|
||||
|
||||
import os
|
||||
import requests
|
||||
|
||||
|
||||
def _short_comment(fn, http_triggers, cron_triggers):
|
||||
"""Генерирует однострочный комментарий-описание функции по её метаданным."""
|
||||
phase = fn.get("phase", "")
|
||||
runtime = fn.get("runtime", "")
|
||||
|
||||
if http_triggers:
|
||||
active_str = "активна" if http_triggers[0].get("active") else "неактивна"
|
||||
return f"HTTP endpoint ({runtime}) — {phase}, {active_str}"
|
||||
elif cron_triggers:
|
||||
schedule = cron_triggers[0].get("schedule", "?")
|
||||
active_str = "активна" if cron_triggers[0].get("active") else "неактивна"
|
||||
return f"Cron '{schedule}' ({runtime}) — {phase}, {active_str}"
|
||||
else:
|
||||
return f"Job/runner без триггера ({runtime}) — {phase}"
|
||||
|
||||
|
||||
def list_all(event):
|
||||
api_url = os.environ["SLESS_API_URL"].rstrip("/")
|
||||
namespace = os.environ["SLESS_NAMESPACE"]
|
||||
token = os.environ["SLESS_TOKEN"]
|
||||
ext_url = os.environ.get("SLESS_EXTERNAL_URL", "").rstrip("/")
|
||||
|
||||
# Имена функций, которые не должны присутствовать в выводе.
|
||||
# Включает саму себя ("funcs") и служебные функции других примеров.
|
||||
exclude = {
|
||||
n.strip()
|
||||
for n in os.environ.get("SLESS_EXCLUDE", "").split(",")
|
||||
if n.strip()
|
||||
}
|
||||
|
||||
headers = {"Authorization": f"Bearer {token}"}
|
||||
|
||||
fns_resp = requests.get(
|
||||
f"{api_url}/v1/namespaces/{namespace}/functions",
|
||||
headers=headers,
|
||||
timeout=10,
|
||||
)
|
||||
fns_resp.raise_for_status()
|
||||
|
||||
trs_resp = requests.get(
|
||||
f"{api_url}/v1/namespaces/{namespace}/triggers",
|
||||
headers=headers,
|
||||
timeout=10,
|
||||
)
|
||||
trs_resp.raise_for_status()
|
||||
|
||||
# Индекс триггеров по имени функции
|
||||
triggers_by_fn = {}
|
||||
for tr in trs_resp.json():
|
||||
fn_name = tr.get("function") or tr.get("functionRef")
|
||||
if fn_name:
|
||||
triggers_by_fn.setdefault(fn_name, []).append(tr)
|
||||
|
||||
result = []
|
||||
for fn in fns_resp.json():
|
||||
name = fn["name"]
|
||||
if name in exclude:
|
||||
continue
|
||||
|
||||
http_triggers = [
|
||||
t for t in triggers_by_fn.get(name, []) if t.get("type") == "http"
|
||||
]
|
||||
cron_triggers = [
|
||||
t for t in triggers_by_fn.get(name, []) if t.get("type") == "cron"
|
||||
]
|
||||
is_active = any(
|
||||
t.get("enabled", True) and t.get("active", False)
|
||||
for t in triggers_by_fn.get(name, [])
|
||||
)
|
||||
|
||||
entry = {
|
||||
# "#" — первый ключ: служит визуальным комментарием при pretty-print
|
||||
"#": _short_comment(fn, http_triggers, cron_triggers),
|
||||
"name": name,
|
||||
"runtime": fn.get("runtime"),
|
||||
"phase": fn.get("phase"),
|
||||
"active": is_active,
|
||||
}
|
||||
|
||||
# URL вычисляем из SLESS_EXTERNAL_URL если задан — state может хранить старый домен
|
||||
if http_triggers:
|
||||
if ext_url:
|
||||
entry["url"] = f"{ext_url}/fn/{namespace}/{name}"
|
||||
else:
|
||||
entry["url"] = http_triggers[0].get("url", "")
|
||||
|
||||
if cron_triggers:
|
||||
entry["cron"] = cron_triggers[0].get("schedule", "")
|
||||
|
||||
if fn.get("message"):
|
||||
entry["message"] = fn["message"]
|
||||
|
||||
# created_at и last_built_at — доступны после обновления оператора до v0.1.32+
|
||||
if fn.get("created_at"):
|
||||
entry["created_at"] = fn["created_at"]
|
||||
if fn.get("last_built_at"):
|
||||
entry["last_built_at"] = fn["last_built_at"]
|
||||
|
||||
result.append(entry)
|
||||
|
||||
# Сортировка: активные вверх, затем по имени
|
||||
result.sort(key=lambda f: (not f["active"], f["name"]))
|
||||
|
||||
return {
|
||||
"namespace": namespace,
|
||||
"count": len(result),
|
||||
"functions": result,
|
||||
}
|
||||
|
||||
@@ -0,0 +1,73 @@
|
||||
|
||||
# resource "nubes_lucee" "app1" {
|
||||
# # Lucee-приложение, зависит от Postgres
|
||||
# resource_name = "lucy_teststand_0"
|
||||
# # resource_realm = "k8s-3.ext.nubes.ru"
|
||||
# resource_realm = nubes_postgres.db2.resource_realm
|
||||
# # resource_realm = "k8s-4-sandbox-nubes-ru"
|
||||
# domain = "web-test-stand"
|
||||
|
||||
# git_path = "https://gitea-naeel.giteak8s.services.ngcloud.ru/naeel/testlucee"
|
||||
|
||||
# json_env = jsonencode({
|
||||
# # 🔗 Настройки Data Source 'testds' для Lucee (Application.cfc)
|
||||
# testds_class = "org.postgresql.Driver" # 📂 Драйвер БД
|
||||
# testds_bundleName = "org.postgresql.jdbc" # 📦 Имя бандла JDBC
|
||||
# testds_bundleVersion = "42.6.0" # 🔢 Версия драйвера
|
||||
# testds_connectionString = "jdbc:postgresql://${nubes_postgres.db2.state_out_flat["internalConnect.master"]}:5432/postgres?sslmode=require" # 🚀 Строка подключения
|
||||
# testds_username = nubes_postgres_user.db2_user.username # 👤 Логин
|
||||
# testds_password = jsondecode(nubes_postgres.db2.vault_secrets["users"])[nubes_postgres_user.db2_user.username]["password"] # 🔑 Пароль
|
||||
# testds_connectionLimit = "5" # 🚦 Лимит соединений
|
||||
# testds_liveTimeout = "15" # ⏳ Таймаут жизни
|
||||
# testds_validate = "false" # ✅ Валидация при запросе
|
||||
# })
|
||||
|
||||
# resource_c_p_u = 300
|
||||
# resource_memory = 512
|
||||
# resource_instances = 1
|
||||
# app_version = "5.4"
|
||||
|
||||
# depends_on = [nubes_postgres.db2]
|
||||
# }
|
||||
|
||||
# resource "nubes_nodejs" "app3" {
|
||||
# # NodeJS демо, работающий с тем же Postgres.
|
||||
# resource_name = "node_01"
|
||||
# resource_realm = nubes_postgres.db2.resource_realm
|
||||
# domain = "node07"
|
||||
# git_path = "https://gitea-naeel.giteak8s.services.ngcloud.ru/naeel/testnode.git"
|
||||
# health_path = "/healthz"
|
||||
# app_version = "23"
|
||||
|
||||
# json_env = jsonencode({
|
||||
# # Переменные подключения к Postgres.
|
||||
# PGHOST = nubes_postgres.db2.state_out_flat["internalConnect.master"]
|
||||
# PGPORT = "5432"
|
||||
# PGUSER = nubes_postgres_user.db2_user.username
|
||||
# PGPASSWORD = jsondecode(nubes_postgres.db2.vault_secrets["users"])[nubes_postgres_user.db2_user.username]["password"]
|
||||
# PGDATABASE = nubes_postgres_database.db2_app.db_name
|
||||
# PGSSLMODE = "require"
|
||||
# DATABASE_URL = format(
|
||||
# "postgresql://%s:%s@%s:5432/%s?sslmode=require",
|
||||
# nubes_postgres_user.db2_user.username,
|
||||
# jsondecode(nubes_postgres.db2.vault_secrets["users"])[nubes_postgres_user.db2_user.username]["password"],
|
||||
# nubes_postgres.db2.state_out_flat["internalConnect.master"],
|
||||
# nubes_postgres_database.db2_app.db_name
|
||||
# )
|
||||
# })
|
||||
|
||||
# resource_c_p_u = 300
|
||||
# resource_memory = 256
|
||||
# resource_instances = 1
|
||||
|
||||
# depends_on = [nubes_postgres.db2]
|
||||
# }
|
||||
|
||||
# output "pg_vault_secrets" {
|
||||
# value = nubes_postgres.db2.vault_secrets
|
||||
# sensitive = true
|
||||
# }
|
||||
|
||||
# terraform output -json pg_vault_secrets
|
||||
|
||||
|
||||
@@ -0,0 +1,58 @@
|
||||
// 2026-03-17 17:05
|
||||
// main.tf — провайдеры и переменные для Nubes + sless.
|
||||
terraform {
|
||||
required_providers {
|
||||
nubes = {
|
||||
source = "terra.k8c.ru/nubes/nubes"
|
||||
version = "5.0.19"
|
||||
}
|
||||
sless = {
|
||||
source = "terra.k8c.ru/naeel/sless"
|
||||
version = "~> 0.1.18"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
variable "api_token" {
|
||||
type = string
|
||||
sensitive = true
|
||||
description = "Nubes API token"
|
||||
}
|
||||
variable "s3_uid" {
|
||||
type = string
|
||||
sensitive = true
|
||||
description = "Nubes S3 UID"
|
||||
}
|
||||
variable "realm" {
|
||||
type = string
|
||||
sensitive = true
|
||||
description = "resource_realm parameter for nubes_postgres resource"
|
||||
}
|
||||
|
||||
// 2026-03-18 — pg_user/pg_password помечены optional (default="") для сверки.
|
||||
// Реальные credentials берутся из vault_secrets через locals в resources.tf.
|
||||
variable "pg_user" {
|
||||
type = string
|
||||
sensitive = true
|
||||
default = ""
|
||||
description = "Только для сверки. Реальный username из nubes_postgres_user.pg_user.username. Должен совпадать с vault."
|
||||
}
|
||||
|
||||
variable "pg_password" {
|
||||
type = string
|
||||
sensitive = true
|
||||
default = ""
|
||||
description = "Только для сверки. Реальный пароль из vault_secrets. Должен совпадать с tfvars."
|
||||
}
|
||||
|
||||
provider "nubes" {
|
||||
api_token = var.api_token
|
||||
api_endpoint = "https://deck-test.ngcloud.ru/api/v1/index.cfm"
|
||||
}
|
||||
|
||||
provider "sless" {
|
||||
endpoint = "https://sless.kube5s.ru"
|
||||
token = var.api_token
|
||||
nubes_endpoint = "https://deck-test.ngcloud.ru/api/v1"
|
||||
}
|
||||
|
||||
@@ -0,0 +1,196 @@
|
||||
// 2026-03-18 — добавлены locals для извлечения credentials из vault_secrets (без хардкода).
|
||||
// Для сверки хардкод остаётся в terraform.tfvars на этапе разработки.
|
||||
// sless_function и sless_job закомментированы — сначала проверяется сетевое соединение.
|
||||
|
||||
# Актуальные credentials из vault_secrets (authoritatively) — vault синхронизирован с кластером.
|
||||
# Структура vault_secrets["users"]: JSON-строка {"username": {"password": "...", "username": "..."}}
|
||||
locals {
|
||||
pg_creds_map = jsondecode(nubes_postgres.npg.vault_secrets["users"])
|
||||
pg_username = nubes_postgres_user.pg_user.username
|
||||
pg_password = local.pg_creds_map[local.pg_username]["password"]
|
||||
pg_host = nubes_postgres.npg.state_out_flat["internalConnect.master"]
|
||||
pg_database = nubes_postgres_database.db.db_name
|
||||
}
|
||||
|
||||
resource "nubes_postgres" "npg" {
|
||||
resource_name = "testnarod-pg-0"
|
||||
# s3_uid = "s01325"
|
||||
s3_uid = var.s3_uid
|
||||
resource_realm = var.realm
|
||||
resource_instances = 1
|
||||
resource_memory = 512
|
||||
resource_c_p_u = 500
|
||||
resource_disk = "1"
|
||||
app_version = "17"
|
||||
json_parameters = jsonencode({
|
||||
log_connections = "off"
|
||||
log_disconnections = "off"
|
||||
})
|
||||
enable_pg_pooler_master = false
|
||||
enable_pg_pooler_slave = false
|
||||
allow_no_s_s_l = false
|
||||
auto_scale = false
|
||||
auto_scale_percentage = 10
|
||||
auto_scale_tech_window = 0
|
||||
auto_scale_quota_gb = "1"
|
||||
need_external_address_master = false
|
||||
|
||||
# suspend_on_destroy = false
|
||||
operation_timeout = "11m"
|
||||
adopt_existing_on_create = true
|
||||
}
|
||||
|
||||
resource "nubes_postgres_user" "pg_user" {
|
||||
postgres_id = nubes_postgres.npg.id
|
||||
username = "u-user0"
|
||||
role = "ddl_user"
|
||||
adopt_existing_on_create = true
|
||||
}
|
||||
|
||||
resource "nubes_postgres_database" "db" {
|
||||
postgres_id = nubes_postgres.npg.id
|
||||
db_name = "db_terra"
|
||||
db_owner = nubes_postgres_user.pg_user.username
|
||||
adopt_existing_on_create = true
|
||||
# suspend_on_destroy = false
|
||||
}
|
||||
|
||||
# Служебная функция выполняет SQL-операторы из event_json.
|
||||
# Credentials берутся из locals (vault_secrets) — без хардкода.
|
||||
# Для сверки хардкод остаётся в terraform.tfvars.
|
||||
resource "sless_function" "postgres_sql_runner_create_table" {
|
||||
name = "pg-create-table-runner"
|
||||
runtime = "python3.11"
|
||||
entrypoint = "sql_runner.run_sql"
|
||||
memory_mb = 128
|
||||
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"
|
||||
# Для сверки (должно совпадать с vault):
|
||||
# PGUSER = var.pg_user
|
||||
# PGPASSWORD = var.pg_password
|
||||
}
|
||||
|
||||
source_dir = "${path.module}/code/sql-runner"
|
||||
}
|
||||
|
||||
resource "sless_job" "postgres_table_init_job" {
|
||||
name = "pg-create-table-job-main-v13"
|
||||
function = sless_function.postgres_sql_runner_create_table.name
|
||||
wait_timeout_sec = 180
|
||||
run_id = 13
|
||||
|
||||
event_json = jsonencode({
|
||||
statements = [
|
||||
"CREATE TABLE IF NOT EXISTS terraform_demo_table (id serial PRIMARY KEY, title text NOT NULL, created_at timestamp DEFAULT now())"
|
||||
]
|
||||
})
|
||||
|
||||
depends_on = [nubes_postgres_database.db]
|
||||
}
|
||||
|
||||
# HTTP-функция на NodeJS: возвращает версию PG-сервера и счётчик строк в таблице.
|
||||
# Единственная функция примера на nodejs20 — проверка что JS runtime работает.
|
||||
# Доступна по URL: https://sless.kube5s.ru/fn/<namespace>/pg-info
|
||||
resource "sless_function" "pg_info" {
|
||||
name = "pg-info"
|
||||
runtime = "nodejs20"
|
||||
entrypoint = "pg_info.info"
|
||||
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/pg-info"
|
||||
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
|
||||
resource "sless_trigger" "pg_info_http" {
|
||||
name = "pg-info-http"
|
||||
type = "http"
|
||||
function = sless_function.pg_info.name
|
||||
enabled = true
|
||||
}
|
||||
|
||||
# HTTP-функции чтения и записи строк terraform_demo_table — в одном файле table_rw.py.
|
||||
# list_rows (GET) — читает все строки; add_row (POST {title}) — вставляет строку.
|
||||
# Доступны по URL: https://sless.kube5s.ru/fn/<namespace>/pg-table-reader
|
||||
# https://sless.kube5s.ru/fn/<namespace>/pg-table-writer
|
||||
resource "sless_function" "postgres_table_reader" {
|
||||
name = "pg-table-reader"
|
||||
runtime = "python3.11"
|
||||
entrypoint = "table_rw.list_rows"
|
||||
memory_mb = 128
|
||||
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/table-rw"
|
||||
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
|
||||
resource "sless_trigger" "postgres_table_reader_http" {
|
||||
name = "pg-table-reader-http"
|
||||
type = "http"
|
||||
function = sless_function.postgres_table_reader.name
|
||||
enabled = true
|
||||
}
|
||||
|
||||
output "table_reader_url" {
|
||||
value = sless_trigger.postgres_table_reader_http.url
|
||||
}
|
||||
|
||||
resource "sless_function" "postgres_table_writer" {
|
||||
name = "pg-table-writer"
|
||||
runtime = "python3.11"
|
||||
entrypoint = "table_rw.add_row"
|
||||
memory_mb = 128
|
||||
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/table-rw"
|
||||
|
||||
depends_on = [sless_job.postgres_table_init_job]
|
||||
}
|
||||
|
||||
resource "sless_trigger" "postgres_table_writer_http" {
|
||||
name = "pg-table-writer-http"
|
||||
type = "http"
|
||||
function = sless_function.postgres_table_writer.name
|
||||
enabled = true
|
||||
}
|
||||
|
||||
output "table_writer_url" {
|
||||
value = sless_trigger.postgres_table_writer_http.url
|
||||
}
|
||||
|
||||
@@ -0,0 +1,51 @@
|
||||
# 2026-03-18 — debug pod для проверки psql-соединения из namespace функций.
|
||||
# Запускается разово. Подключается к тому же postgres, что и sless_function.
|
||||
# kubectl apply -f /tmp/pg-debug-pod.yaml
|
||||
# kubectl logs -n sless-fn-sless-ffd1f598c169b0ae pg-debug-pod
|
||||
|
||||
apiVersion: v1
|
||||
kind: Pod
|
||||
metadata:
|
||||
name: pg-debug-pod
|
||||
namespace: sless-fn-sless-ffd1f598c169b0ae
|
||||
labels:
|
||||
purpose: debug-postgres-connectivity
|
||||
spec:
|
||||
restartPolicy: Never
|
||||
containers:
|
||||
- name: psql
|
||||
image: postgres:17-alpine
|
||||
command:
|
||||
- sh
|
||||
- -c
|
||||
- |
|
||||
echo "=== Testing TCP connectivity to postgres ==="
|
||||
nc -zv -w5 $PGHOST 5432 && echo "TCP OK" || echo "TCP FAILED"
|
||||
|
||||
echo ""
|
||||
echo "=== Testing psql connection ==="
|
||||
PGCONNECT_TIMEOUT=10 psql \
|
||||
"host=$PGHOST port=$PGPORT dbname=$PGDATABASE user=$PGUSER sslmode=$PGSSLMODE" \
|
||||
--command="SELECT current_user, current_database(), version();" \
|
||||
2>&1
|
||||
|
||||
echo ""
|
||||
echo "=== Listing tables ==="
|
||||
PGCONNECT_TIMEOUT=10 psql \
|
||||
"host=$PGHOST port=$PGPORT dbname=$PGDATABASE user=$PGUSER sslmode=$PGSSLMODE" \
|
||||
--command="\dt" \
|
||||
2>&1
|
||||
env:
|
||||
- name: PGHOST
|
||||
value: "postgresqlk8s-master.36875359-dcea-48c4-a593-b4531f20fe96.svc.cluster.local"
|
||||
- name: PGPORT
|
||||
value: "5432"
|
||||
- name: PGDATABASE
|
||||
value: "db_terra"
|
||||
- name: PGUSER
|
||||
value: "u-user0"
|
||||
- name: PGPASSWORD
|
||||
# Актуальный пароль из vault_secrets (совпадает с tfvars.pg_password на 2026-03-18)
|
||||
value: "M03O6fRsngWcVHB2YGivyLfbfxoii2R21nyh2A2r7WSZS5deLwBgLKkc9Wk24Zyl"
|
||||
- name: PGSSLMODE
|
||||
value: "require"
|
||||
@@ -0,0 +1,40 @@
|
||||
# 2026-03-17 13:05
|
||||
# read_pg_user_secret.py — читает пароль пользователя managed PostgreSQL из k8s Secret.
|
||||
# Используется из Terraform external data source, чтобы apply сам получал актуальный пароль
|
||||
# даже для уже существующего пользователя, созданного вне текущего state.
|
||||
|
||||
import base64
|
||||
import json
|
||||
import subprocess
|
||||
import sys
|
||||
|
||||
|
||||
def main():
|
||||
# Читаем query от Terraform external provider из stdin.
|
||||
query = json.load(sys.stdin)
|
||||
namespace = query["namespace"]
|
||||
secret_name = query["secret"]
|
||||
|
||||
# kubectl уже настроен на удалённой машине; читаем ровно поле data.password.
|
||||
result = subprocess.run(
|
||||
[
|
||||
"kubectl",
|
||||
"get",
|
||||
"secret",
|
||||
"-n",
|
||||
namespace,
|
||||
secret_name,
|
||||
"-o",
|
||||
"jsonpath={.data.password}",
|
||||
],
|
||||
check=True,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
)
|
||||
|
||||
password = base64.b64decode(result.stdout.strip()).decode()
|
||||
json.dump({"password": password}, sys.stdout)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -16,7 +16,7 @@ variable "sless_endpoint" {
|
||||
variable "nubes_endpoint" {
|
||||
description = "Nubes endpoint (нужен провайдеру)"
|
||||
type = string
|
||||
default = "https://deck-api-test.ngcloud.ru/api/v1"
|
||||
default = "https://deck-test.ngcloud.ru/api/v1"
|
||||
}
|
||||
|
||||
variable "pg_dsn" {
|
||||
|
||||
@@ -13,5 +13,5 @@ terraform {
|
||||
provider "sless" {
|
||||
endpoint = "https://sless.kube5s.ru"
|
||||
token = var.token
|
||||
nubes_endpoint = "https://deck-api-test.ngcloud.ru/api/v1"
|
||||
nubes_endpoint = "https://deck-test.ngcloud.ru/api/v1"
|
||||
}
|
||||
|
||||
@@ -19,6 +19,6 @@ terraform {
|
||||
provider "sless" {
|
||||
endpoint = "https://sless.kube5s.ru"
|
||||
token = var.token
|
||||
nubes_endpoint = "https://deck-api-test.ngcloud.ru/api/v1"
|
||||
nubes_endpoint = "https://deck-test.ngcloud.ru/api/v1"
|
||||
}
|
||||
|
||||
|
||||
@@ -23,5 +23,5 @@ terraform {
|
||||
provider "sless" {
|
||||
endpoint = "https://sless.kube5s.ru"
|
||||
token = var.token
|
||||
nubes_endpoint = "https://deck-api-test.ngcloud.ru/api/v1"
|
||||
nubes_endpoint = "https://deck-test.ngcloud.ru/api/v1"
|
||||
}
|
||||
|
||||
@@ -13,5 +13,5 @@ terraform {
|
||||
provider "sless" {
|
||||
endpoint = "https://sless.kube5s.ru"
|
||||
token = var.token
|
||||
nubes_endpoint = "https://deck-api-test.ngcloud.ru/api/v1"
|
||||
nubes_endpoint = "https://deck-test.ngcloud.ru/api/v1"
|
||||
}
|
||||
|
||||
@@ -27,5 +27,5 @@ terraform {
|
||||
provider "sless" {
|
||||
endpoint = "https://sless.kube5s.ru"
|
||||
token = var.token
|
||||
nubes_endpoint = "https://deck-api-test.ngcloud.ru/api/v1"
|
||||
nubes_endpoint = "https://deck-test.ngcloud.ru/api/v1"
|
||||
}
|
||||
|
||||
@@ -26,5 +26,5 @@ terraform {
|
||||
provider "sless" {
|
||||
endpoint = "https://sless.kube5s.ru"
|
||||
token = var.token
|
||||
nubes_endpoint = "https://deck-api-test.ngcloud.ru/api/v1"
|
||||
nubes_endpoint = "https://deck-test.ngcloud.ru/api/v1"
|
||||
}
|
||||
|
||||
@@ -54,6 +54,7 @@ require (
|
||||
github.com/prometheus/client_model v0.3.0 // indirect
|
||||
github.com/prometheus/common v0.37.0 // indirect
|
||||
github.com/prometheus/procfs v0.8.0 // indirect
|
||||
github.com/rabbitmq/amqp091-go v1.10.0 // indirect
|
||||
github.com/rs/xid v1.6.0 // indirect
|
||||
github.com/spf13/pflag v1.0.5 // indirect
|
||||
github.com/tinylib/msgp v1.6.1 // indirect
|
||||
|
||||
@@ -270,6 +270,8 @@ github.com/prometheus/procfs v0.6.0/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1
|
||||
github.com/prometheus/procfs v0.7.3/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1xBZuNvfVA=
|
||||
github.com/prometheus/procfs v0.8.0 h1:ODq8ZFEaYeCaZOJlZZdJA2AbQR98dSHSM1KW/You5mo=
|
||||
github.com/prometheus/procfs v0.8.0/go.mod h1:z7EfXMXOkbkqb9IINtpCn86r/to3BnA0uaxHdg830/4=
|
||||
github.com/rabbitmq/amqp091-go v1.10.0 h1:STpn5XsHlHGcecLmMFCtg7mqq0RnD+zFr4uzukfVhBw=
|
||||
github.com/rabbitmq/amqp091-go v1.10.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o=
|
||||
github.com/rogpeppe/go-internal v1.3.0/go.mod h1:M8bDsm7K2OlrFYOpmOWEs/qY81heoFRclV5y23lUDJ4=
|
||||
github.com/rs/xid v1.6.0 h1:fV591PaemRlL6JfRxGDEPl69wICngIQ3shQtzfy2gxU=
|
||||
github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0=
|
||||
@@ -305,6 +307,7 @@ go.uber.org/atomic v1.7.0/go.mod h1:fEN4uk6kAWBTFdckzkM89CLk9XfWZrxpCo0nPH17wJc=
|
||||
go.uber.org/goleak v1.1.10/go.mod h1:8a7PlsEVH3e/a/GLqe5IIrQx6GzcnRmZEufDUTk4A7A=
|
||||
go.uber.org/goleak v1.2.0 h1:xqgm/S+aQvhWFTtR0XK3Jvg7z8kGV8P4X14IzwN3Eqk=
|
||||
go.uber.org/goleak v1.2.0/go.mod h1:XJYK+MuIchqpmGmUSAzotztawfKvYLUIgg7guXrwVUo=
|
||||
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
|
||||
go.uber.org/multierr v1.6.0 h1:y6IPFStTAIT5Ytl7/XYmHvzXQ7S3g/IeZW9hyZ5thw4=
|
||||
go.uber.org/multierr v1.6.0/go.mod h1:cdWPpRnG4AhwMwsgIHip0KRBQjJy5kYEpYjJxpXp9iU=
|
||||
go.uber.org/zap v1.19.0/go.mod h1:xg/QME4nWcxGxrpdeYfq7UvYrLh66cuVKdrbD1XF/NI=
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
// Изменено: 2026-03-11
|
||||
// Изменено: 2026-03-12
|
||||
// invoke.go — прокси-обработчик для вызова HTTP-триггеров функций.
|
||||
// Маршрут: ANY /fn/{namespace}/{name} и /fn/{namespace}/{name}/**
|
||||
// Не защищён auth-токеном — это публичный эндпоинт для вызова функций.
|
||||
@@ -6,99 +6,120 @@
|
||||
// http://{name}.sless-fn-{namespace}.svc.cluster.local:8080
|
||||
// Sub-path и query string пробрасываются как есть:
|
||||
// /fn/ns/notes/add?title=x → http://notes.sless-fn-ns.svc.../add?title=x
|
||||
// Таймаут берётся из Spec.TimeoutSec функции (+ 5s буфер) чтобы не резать
|
||||
// длительные вызовы (stress-тесты, batch-задачи).
|
||||
|
||||
package handler
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/gorilla/mux"
|
||||
slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
)
|
||||
|
||||
// httpClient используется для обращения к функциям внутри кластера.
|
||||
// Таймаут 30s — достаточно для холодного старта функции.
|
||||
var httpClient = &http.Client{Timeout: 30 * time.Second}
|
||||
|
||||
// hopByHopHeaders — заголовки которые нельзя пробрасывать через прокси (RFC 2616 §13.5.1).
|
||||
// Они управляют соединением между двумя узлами, а не end-to-end.
|
||||
// Особо опасен Transfer-Encoding: если пробросить его, клиент неверно интерпретирует тело.
|
||||
var hopByHopHeaders = map[string]bool{
|
||||
"Connection": true,
|
||||
"Keep-Alive": true,
|
||||
"Proxy-Authenticate": true,
|
||||
"Proxy-Authorization": true,
|
||||
"Te": true,
|
||||
"Trailers": true,
|
||||
"Transfer-Encoding": true,
|
||||
"Upgrade": true,
|
||||
"Connection": true,
|
||||
"Keep-Alive": true,
|
||||
"Proxy-Authenticate": true,
|
||||
"Proxy-Authorization": true,
|
||||
"Te": true,
|
||||
"Trailers": true,
|
||||
"Transfer-Encoding": true,
|
||||
"Upgrade": true,
|
||||
}
|
||||
|
||||
// invokeHTTPClient создаёт http.Client с таймаутом под конкретный вызов.
|
||||
// timeout = TimeoutSec функции + 5s буфер на сетевые задержки.
|
||||
// Если TimeoutSec == 0 (не задан), используем 30s по умолчанию.
|
||||
func invokeHTTPClient(timeoutSec int32) *http.Client {
|
||||
t := time.Duration(timeoutSec)*time.Second + 5*time.Second
|
||||
if timeoutSec <= 0 {
|
||||
t = 30 * time.Second
|
||||
}
|
||||
return &http.Client{Timeout: t}
|
||||
}
|
||||
|
||||
// InvokeFunction проксирует входящий запрос к Service функции в кластере.
|
||||
// Namespace выбирается из пути, имя функции — тоже из пути.
|
||||
// Сохраняет метод, тело, Content-Type, sub-path и query string.
|
||||
// Таймаут прокси-клиента = Spec.TimeoutSec функции + 5s (резинка).
|
||||
func (h *Handler) InvokeFunction(w http.ResponseWriter, r *http.Request) {
|
||||
vars := mux.Vars(r)
|
||||
ns := vars["namespace"]
|
||||
name := vars["name"]
|
||||
vars := mux.Vars(r)
|
||||
ns := vars["namespace"]
|
||||
name := vars["name"]
|
||||
|
||||
// Вычисляем sub-path после /fn/{namespace}/{name}
|
||||
// Например: /fn/default/notes/add → subPath = /add
|
||||
prefix := "/fn/" + ns + "/" + name
|
||||
subPath := strings.TrimPrefix(r.URL.Path, prefix)
|
||||
|
||||
// Внутренний URL к Service функции (DNS внутри кластера)
|
||||
target := fmt.Sprintf("http://%s.sless-fn-%s.svc.cluster.local:8080%s", name, ns, subPath)
|
||||
|
||||
// Пробрасываем query string если есть
|
||||
if r.URL.RawQuery != "" {
|
||||
target += "?" + r.URL.RawQuery
|
||||
}
|
||||
|
||||
// Создаём проксируемый запрос с тем же методом и телом
|
||||
proxyReq, err := http.NewRequestWithContext(r.Context(), r.Method, target, r.Body)
|
||||
if err != nil {
|
||||
h.Log.Error("invoke: failed to create proxy request", "err", err, "ns", ns, "fn", name)
|
||||
writeJSON(w, http.StatusInternalServerError, errResp("failed to create proxy request"))
|
||||
return
|
||||
}
|
||||
|
||||
// Пробрасываем Content-Type и Content-Length если есть.
|
||||
// Content-Length обязателен: Python BaseHTTPRequestHandler читает тело
|
||||
// ровно столько байт, сколько указано в заголовке; без него body = пусто.
|
||||
if ct := r.Header.Get("Content-Type"); ct != "" {
|
||||
proxyReq.Header.Set("Content-Type", ct)
|
||||
}
|
||||
proxyReq.ContentLength = r.ContentLength
|
||||
|
||||
resp, err := httpClient.Do(proxyReq)
|
||||
if err != nil {
|
||||
// "no such host" — Service не существует (функция удалена), возвращаем 404.
|
||||
// Это отличается от временной сетевой ошибки: NXDOMAIN строго означает отсутствие записи.
|
||||
if strings.Contains(err.Error(), "no such host") {
|
||||
writeJSON(w, http.StatusNotFound, errResp("function not found"))
|
||||
return
|
||||
}
|
||||
h.Log.Error("invoke: function unreachable", "err", err, "ns", ns, "fn", name, "target", target)
|
||||
writeJSON(w, http.StatusBadGateway, errResp("function unreachable: "+err.Error()))
|
||||
return
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
// Копируем заголовки и статус из ответа функции.
|
||||
// Hop-by-hop заголовки фильтруем: они управляют конкретным TCP-соединением
|
||||
// и не должны пробрасываться через прокси (RFC 2616 §13.5.1).
|
||||
for k, vals := range resp.Header {
|
||||
if hopByHopHeaders[k] {
|
||||
continue
|
||||
}
|
||||
for _, v := range vals {
|
||||
w.Header().Add(k, v)
|
||||
}
|
||||
}
|
||||
w.WriteHeader(resp.StatusCode)
|
||||
_, _ = io.Copy(w, resp.Body)
|
||||
// Смотрим TimeoutSec из Function CRD, чтобы не резать длительные вызовы.
|
||||
// Если функция не найдена — продолжаем с дефолтным таймаутом (30s).
|
||||
var timeoutSec int32
|
||||
fn := &slessv1alpha1.Function{}
|
||||
if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, fn); err == nil {
|
||||
timeoutSec = fn.Spec.TimeoutSec
|
||||
}
|
||||
httpClient := invokeHTTPClient(timeoutSec)
|
||||
|
||||
// Вычисляем sub-path после /fn/{namespace}/{name}
|
||||
// Например: /fn/default/notes/add → subPath = /add
|
||||
prefix := "/fn/" + ns + "/" + name
|
||||
subPath := strings.TrimPrefix(r.URL.Path, prefix)
|
||||
|
||||
// Внутренний URL к Service функции (DNS внутри кластера)
|
||||
target := fmt.Sprintf("http://%s.sless-fn-%s.svc.cluster.local:8080%s", name, ns, subPath)
|
||||
|
||||
// Пробрасываем query string если есть
|
||||
if r.URL.RawQuery != "" {
|
||||
target += "?" + r.URL.RawQuery
|
||||
}
|
||||
|
||||
// Создаём проксируемый запрос с тем же методом и телом
|
||||
proxyReq, err := http.NewRequestWithContext(r.Context(), r.Method, target, r.Body)
|
||||
if err != nil {
|
||||
h.Log.Error("invoke: failed to create proxy request", "err", err, "ns", ns, "fn", name)
|
||||
writeJSON(w, http.StatusInternalServerError, errResp("failed to create proxy request"))
|
||||
return
|
||||
}
|
||||
|
||||
// Пробрасываем Content-Type и Content-Length если есть.
|
||||
// Content-Length обязателен: Python BaseHTTPRequestHandler читает тело
|
||||
// ровно столько байт, сколько указано в заголовке; без него body = пусто.
|
||||
if ct := r.Header.Get("Content-Type"); ct != "" {
|
||||
proxyReq.Header.Set("Content-Type", ct)
|
||||
}
|
||||
proxyReq.ContentLength = r.ContentLength
|
||||
|
||||
resp, err := httpClient.Do(proxyReq)
|
||||
if err != nil {
|
||||
// "no such host" — Service не существует (функция удалена), возвращаем 404.
|
||||
// Это отличается от временной сетевой ошибки: NXDOMAIN строго означает отсутствие записи.
|
||||
if strings.Contains(err.Error(), "no such host") {
|
||||
writeJSON(w, http.StatusNotFound, errResp("function not found"))
|
||||
return
|
||||
}
|
||||
h.Log.Error("invoke: function unreachable", "err", err, "ns", ns, "fn", name, "target", target)
|
||||
writeJSON(w, http.StatusBadGateway, errResp("function unreachable: "+err.Error()))
|
||||
return
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
// Копируем заголовки и статус из ответа функции.
|
||||
// Hop-by-hop заголовки фильтруем: они управляют конкретным TCP-соединением
|
||||
// и не должны пробрасываться через прокси (RFC 2616 §13.5.1).
|
||||
for k, vals := range resp.Header {
|
||||
if hopByHopHeaders[k] {
|
||||
continue
|
||||
}
|
||||
for _, v := range vals {
|
||||
w.Header().Add(k, v)
|
||||
}
|
||||
}
|
||||
w.WriteHeader(resp.StatusCode)
|
||||
_, _ = io.Copy(w, resp.Body)
|
||||
}
|
||||
|
||||
@@ -125,7 +125,7 @@ func (b *Builder) Build(ctx context.Context, namespace, funcName, s3Key string)
|
||||
Args: []string{
|
||||
"--context=" + s3ContextURL,
|
||||
"--destination=" + imageRef,
|
||||
"--no-cache", // без кэша — гарантирует что COPY берёт свежий код из S3
|
||||
// --no-cache не поддерживается этой версией kaniko; кэш отключён по умолчанию
|
||||
},
|
||||
Env: []corev1.EnvVar{
|
||||
{Name: "AWS_ACCESS_KEY_ID", Value: b.s3AccessKey},
|
||||
|
||||
@@ -70,9 +70,9 @@ func runtimeBaseImage(runtime string) (string, error) {
|
||||
case "nodejs20":
|
||||
return "naeel/sless-runtime-nodejs20:v0.1.2", nil
|
||||
case "go1.23":
|
||||
// Go использует builder-образ (содержит golang:1.23 + server.go).
|
||||
// Go использует builder-образ (содержит golang:1.23 + server.go + go.mod/go.sum с pgx/v5).
|
||||
// Конечный образ multi-stage — alpine без компилятора.
|
||||
return "naeel/sless-runtime-go1.23:v0.1.0", nil
|
||||
return "naeel/sless-runtime-go1.23:v0.1.1", nil
|
||||
default:
|
||||
return "", fmt.Errorf("unsupported runtime: %q (supported: python3.11, nodejs20, go1.23)", runtime)
|
||||
}
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
// Изменено: 2026-03-11
|
||||
// Изменено: 2026-03-19
|
||||
// Конфигурация сервиса — читается из env переменных при старте.
|
||||
// Все компоненты (API, builder, runner) получают конфиг через эту структуру.
|
||||
// Используем env а не файлы конфигурации — стандарт для k8s (ConfigMap/Secret → env).
|
||||
@@ -56,6 +56,11 @@ type Config struct {
|
||||
// Если не заданы, EnsureProject пропускается (только DockerHub/без проектов).
|
||||
HarborUser string
|
||||
HarborPass string
|
||||
|
||||
// RabbitMQURL — AMQP URL брокера для event-dispatcher.
|
||||
// Формат: amqp://user:pass@host:5672/
|
||||
// Опционально — если не задан, event-триггеры не будут обрабатываться.
|
||||
RabbitMQURL string
|
||||
}
|
||||
|
||||
func Load() (*Config, error) {
|
||||
@@ -133,9 +138,12 @@ func Load() (*Config, error) {
|
||||
cfg.APIToken = os.Getenv("SLESS_API_TOKEN")
|
||||
|
||||
// Harbor API креды — опциональны.
|
||||
// Если не заданы, EnsureProject пропускается (push работает через docker-кред в kaniko).
|
||||
// Если не заданы, EnsureProject пропускается (push работает через docker-кред в канiko).
|
||||
cfg.HarborUser = os.Getenv("HARBOR_USER")
|
||||
cfg.HarborPass = os.Getenv("HARBOR_PASS")
|
||||
|
||||
// RabbitMQ — опционально. Нужен только если используются event-триггеры.
|
||||
cfg.RabbitMQURL = os.Getenv("RABBITMQ_URL")
|
||||
|
||||
return cfg, nil
|
||||
}
|
||||
|
||||
+19
-23
@@ -1,33 +1,29 @@
|
||||
# Создано: 2026-03-11
|
||||
# Изменено: 2026-03-19 — v0.1.1: добавлены go.sum + go mod download (pgx/v5 v5.7.2).
|
||||
# Base builder image для Go 1.23 serverless функций.
|
||||
# Это BUILDER-образ (не runtime): содержит Go компилятор + server.go.
|
||||
# server.go — HTTP-wrapper + job-runner, компилируется ВМЕСТЕ с кодом пользователя.
|
||||
# Это BUILDER-образ: golang:1.23-alpine + server.go + pre-cached зависимости.
|
||||
# go mod download кеширует pgx/v5 в /root/go/pkg/mod — kaniko не скачивает их при каждой сборке.
|
||||
#
|
||||
# kaniko делает multi-stage build:
|
||||
# Stage 1: golang:1.23-alpine → COPY server.go + handler/ → go build → /server
|
||||
# Stage 2: FROM alpine → COPY /server → минимальный финальный образ.
|
||||
# kaniko генерирует Dockerfile для каждой функции:
|
||||
# FROM naeel/sless-runtime-go1.23:v0.1.1 AS builder
|
||||
# WORKDIR /app
|
||||
# COPY . /app/handler/
|
||||
# RUN CGO_ENABLED=0 go build -o /server .
|
||||
# FROM alpine:3.20
|
||||
# COPY --from=builder /server /server → финальный образ
|
||||
#
|
||||
# Почему builder-образ, а не runtime: Go требует статической компиляции — нельзя
|
||||
# "подключить" пользовательский код динамически как в Python/Node.
|
||||
# Почему не multi-stage здесь: base образ должен иметь Go компилятор и module cache.
|
||||
# Бинарник собирается kaniko, а не здесь.
|
||||
|
||||
FROM golang:1.23-alpine AS builder
|
||||
FROM golang:1.23-alpine
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
# server.go — точка входа (main package), вызывает Handler() из пользовательского кода.
|
||||
# Копируем mod файлы первыми — отдельный кешируемый слой.
|
||||
# При изменении go.mod/go.sum слой инвалидируется и go mod download запускается заново.
|
||||
COPY go.mod go.sum /app/
|
||||
RUN go mod download
|
||||
|
||||
# server.go — HTTP-wrapper + job-runner (main package, импортирует sless/fn/handler).
|
||||
COPY server.go /app/server.go
|
||||
|
||||
# Код пользователя kaniko копирует в /app/handler/
|
||||
COPY handler/ /app/handler/
|
||||
|
||||
# Статическая сборка (без libc) — работает в alpine-контейнере
|
||||
RUN CGO_ENABLED=0 go build -o /server /app/server.go
|
||||
|
||||
# Финальный минимальный образ
|
||||
FROM alpine:3.20
|
||||
|
||||
COPY --from=builder /server /server
|
||||
|
||||
EXPOSE 8080
|
||||
|
||||
CMD ["/server"]
|
||||
|
||||
@@ -1,3 +1,14 @@
|
||||
module sless/fn
|
||||
|
||||
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,197 @@
|
||||
// Изменено: 2026-03-19
|
||||
// Dispatcher — управляет AMQP consumer-ами.
|
||||
// Для каждого Trigger{type:event} держит один goroutine-consumer:
|
||||
// - Subscribe(key, queue, targetURL) — открывает consumer на очередь
|
||||
// - Unsubscribe(key) — закрывает consumer
|
||||
//
|
||||
// При получении сообщения из очереди: POST на targetURL с телом сообщения.
|
||||
// 2xx → ack, остальное → nack (requeue=true).
|
||||
// При потере AMQP соединения — переподключение с экспоненциальной задержкой.
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/go-logr/logr"
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// consumerEntry — запись об активном consumer-е.
|
||||
type consumerEntry struct {
|
||||
queue string
|
||||
targetURL string
|
||||
cancel context.CancelFunc
|
||||
}
|
||||
|
||||
// Dispatcher управляет пулом AMQP consumer-ов (один на каждый event Trigger).
|
||||
type Dispatcher struct {
|
||||
rabbitURL string
|
||||
logger logr.Logger
|
||||
httpClient *http.Client
|
||||
|
||||
mu sync.Mutex
|
||||
consumers map[string]*consumerEntry // key = namespace/name триггера
|
||||
}
|
||||
|
||||
// NewDispatcher создаёт Dispatcher.
|
||||
func NewDispatcher(rabbitURL string, logger logr.Logger) *Dispatcher {
|
||||
return &Dispatcher{
|
||||
rabbitURL: rabbitURL,
|
||||
logger: logger.WithName("dispatcher"),
|
||||
httpClient: &http.Client{Timeout: 30 * time.Second},
|
||||
consumers: make(map[string]*consumerEntry),
|
||||
}
|
||||
}
|
||||
|
||||
// Subscribe регистрирует consumer на очередь queue.
|
||||
// key — уникальный идентификатор триггера (namespace/name).
|
||||
// targetURL — внутренний HTTP URL функции (http://svc.ns.svc.cluster.local:8080/).
|
||||
// Если consumer для этого key уже есть — он переподписывается с новыми параметрами.
|
||||
func (d *Dispatcher) Subscribe(key, queue, targetURL string) {
|
||||
d.mu.Lock()
|
||||
defer d.mu.Unlock()
|
||||
|
||||
// Если уже подписан — отменяем старый и создаём новый
|
||||
if entry, ok := d.consumers[key]; ok {
|
||||
if entry.queue == queue && entry.targetURL == targetURL {
|
||||
return // ничего не изменилось
|
||||
}
|
||||
entry.cancel()
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
entry := &consumerEntry{queue: queue, targetURL: targetURL, cancel: cancel}
|
||||
d.consumers[key] = entry
|
||||
|
||||
go d.runConsumer(ctx, key, queue, targetURL)
|
||||
d.logger.Info("subscribed", "key", key, "queue", queue, "target", targetURL)
|
||||
}
|
||||
|
||||
// Unsubscribe закрывает consumer для триггера key.
|
||||
func (d *Dispatcher) Unsubscribe(key string) {
|
||||
d.mu.Lock()
|
||||
defer d.mu.Unlock()
|
||||
|
||||
if entry, ok := d.consumers[key]; ok {
|
||||
entry.cancel()
|
||||
delete(d.consumers, key)
|
||||
d.logger.Info("unsubscribed", "key", key)
|
||||
}
|
||||
}
|
||||
|
||||
// runConsumer — горутина одного consumer-а.
|
||||
// Держит AMQP соединение и channel. При разрыве переподключается.
|
||||
// Завершается когда ctx отменён.
|
||||
func (d *Dispatcher) runConsumer(ctx context.Context, key, queue, targetURL string) {
|
||||
backoff := time.Second
|
||||
for {
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
err := d.consumeLoop(ctx, queue, targetURL)
|
||||
if ctx.Err() != nil {
|
||||
return // нормальное завершение
|
||||
}
|
||||
d.logger.Error(err, "consumer loop error, reconnecting", "key", key, "backoff", backoff)
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-time.After(backoff):
|
||||
}
|
||||
// Экспоненциальная задержка, максимум 30 секунд
|
||||
if backoff < 30*time.Second {
|
||||
backoff *= 2
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// consumeLoop устанавливает соединение, объявляет очередь и читает сообщения.
|
||||
// Возвращает ошибку при потере соединения (вызывающий перезапустит).
|
||||
func (d *Dispatcher) consumeLoop(ctx context.Context, queue, targetURL string) error {
|
||||
conn, err := amqp.Dial(d.rabbitURL)
|
||||
if err != nil {
|
||||
return fmt.Errorf("amqp dial: %w", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
ch, err := conn.Channel()
|
||||
if err != nil {
|
||||
return fmt.Errorf("amqp channel: %w", err)
|
||||
}
|
||||
defer ch.Close()
|
||||
|
||||
// Объявляем очередь как durable — она переживёт рестарт брокера
|
||||
_, err = ch.QueueDeclare(queue, true, false, false, false, nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("queue declare %q: %w", queue, err)
|
||||
}
|
||||
|
||||
// prefetch=1: не берём следующее сообщение пока не обработали текущее
|
||||
if err = ch.Qos(1, 0, false); err != nil {
|
||||
return fmt.Errorf("qos: %w", err)
|
||||
}
|
||||
|
||||
msgs, err := ch.Consume(queue, "" /*auto-tag*/, false /*autoAck*/, false, false, false, nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("consume: %w", err)
|
||||
}
|
||||
|
||||
connClose := conn.NotifyClose(make(chan *amqp.Error, 1))
|
||||
d.logger.Info("consuming", "queue", queue, "target", targetURL)
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return nil
|
||||
case amqpErr := <-connClose:
|
||||
return fmt.Errorf("connection closed: %v", amqpErr)
|
||||
case msg, ok := <-msgs:
|
||||
if !ok {
|
||||
return fmt.Errorf("messages channel closed")
|
||||
}
|
||||
d.handleMessage(ctx, msg, targetURL)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// handleMessage выполняет HTTP POST на targetURL, ack/nack по результату.
|
||||
func (d *Dispatcher) handleMessage(ctx context.Context, msg amqp.Delivery, targetURL string) {
|
||||
logger := d.logger.WithValues("target", targetURL, "contentType", msg.ContentType)
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, targetURL, bytes.NewReader(msg.Body))
|
||||
if err != nil {
|
||||
logger.Error(err, "failed to build request, nacking")
|
||||
_ = msg.Nack(false, true)
|
||||
return
|
||||
}
|
||||
|
||||
contentType := msg.ContentType
|
||||
if contentType == "" {
|
||||
contentType = "application/octet-stream"
|
||||
}
|
||||
req.Header.Set("Content-Type", contentType)
|
||||
req.ContentLength = int64(len(msg.Body))
|
||||
|
||||
resp, err := d.httpClient.Do(req)
|
||||
if err != nil {
|
||||
logger.Error(err, "http call failed, nacking")
|
||||
_ = msg.Nack(false, true)
|
||||
return
|
||||
}
|
||||
resp.Body.Close()
|
||||
|
||||
if resp.StatusCode >= 200 && resp.StatusCode < 300 {
|
||||
_ = msg.Ack(false)
|
||||
logger.V(1).Info("message delivered", "status", resp.StatusCode)
|
||||
} else {
|
||||
// Функция вернула ошибку — requeue для повтора
|
||||
logger.Info("function returned non-2xx, nacking", "status", resp.StatusCode)
|
||||
_ = msg.Nack(false, true)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,59 @@
|
||||
// Package main — точка входа event-dispatcher.
|
||||
// Изменено: 2026-03-19
|
||||
// Читает RABBITMQ_URL из env, запускает Watcher (следит за Trigger CRD)
|
||||
// и Dispatcher (держит AMQP consumer-ы, роутит сообщения в функции).
|
||||
// Graceful shutdown по SIGTERM/SIGINT.
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
|
||||
"k8s.io/client-go/rest"
|
||||
"k8s.io/client-go/tools/clientcmd"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/log/zap"
|
||||
)
|
||||
|
||||
func main() {
|
||||
ctrl.SetLogger(zap.New(zap.UseDevMode(false)))
|
||||
logger := ctrl.Log.WithName("event-dispatcher")
|
||||
|
||||
rabbitURL := os.Getenv("RABBITMQ_URL")
|
||||
if rabbitURL == "" {
|
||||
logger.Error(nil, "RABBITMQ_URL is required")
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
// k8s config: в кластере — InClusterConfig, снаружи — KUBECONFIG
|
||||
cfg, err := rest.InClusterConfig()
|
||||
if err != nil {
|
||||
// fallback для локальной разработки/тестов
|
||||
kubeconfig := os.Getenv("KUBECONFIG")
|
||||
cfg, err = clientcmd.BuildConfigFromFlags("", kubeconfig)
|
||||
if err != nil {
|
||||
logger.Error(err, "failed to build k8s config")
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
|
||||
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
|
||||
defer cancel()
|
||||
|
||||
disp := NewDispatcher(rabbitURL, logger)
|
||||
watcher, err := NewWatcher(cfg, disp, logger)
|
||||
if err != nil {
|
||||
logger.Error(err, "failed to create watcher")
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
logger.Info("starting event-dispatcher")
|
||||
if err := watcher.Run(ctx); err != nil {
|
||||
logger.Error(err, "watcher stopped with error")
|
||||
os.Exit(1)
|
||||
}
|
||||
logger.Info("event-dispatcher stopped")
|
||||
}
|
||||
@@ -0,0 +1,129 @@
|
||||
// Изменено: 2026-03-19
|
||||
// Watcher — следит за Trigger CRD через k8s informer.
|
||||
// При появлении Trigger{type:event} вызывает Dispatcher.Subscribe.
|
||||
// При удалении — Dispatcher.Unsubscribe.
|
||||
// При изменении (другая очередь или функция) — Subscribe автоматически переподключает.
|
||||
//
|
||||
// targetURL строится как:
|
||||
// http://{functionRef}.{namespace}.svc.cluster.local:8080/
|
||||
// Это внутренний адрес Service функции — оператор создаёт его при reconcileEvent.
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"github.com/go-logr/logr"
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
|
||||
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
|
||||
"k8s.io/client-go/rest"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
|
||||
slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1"
|
||||
)
|
||||
|
||||
var watcherScheme = runtime.NewScheme()
|
||||
|
||||
func init() {
|
||||
utilruntime.Must(clientgoscheme.AddToScheme(watcherScheme))
|
||||
utilruntime.Must(slessv1alpha1.AddToScheme(watcherScheme))
|
||||
}
|
||||
|
||||
// Watcher запускает controller-runtime manager который следит за Trigger CRD.
|
||||
type Watcher struct {
|
||||
manager ctrl.Manager
|
||||
dispatcher *Dispatcher
|
||||
logger logr.Logger
|
||||
}
|
||||
|
||||
// NewWatcher создаёт Watcher.
|
||||
func NewWatcher(cfg *rest.Config, disp *Dispatcher, logger logr.Logger) (*Watcher, error) {
|
||||
mgr, err := ctrl.NewManager(cfg, ctrl.Options{
|
||||
Scheme: watcherScheme,
|
||||
MetricsBindAddress: "0", // метрики не нужны — это не оператор
|
||||
LeaderElection: false,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create manager: %w", err)
|
||||
}
|
||||
|
||||
w := &Watcher{
|
||||
manager: mgr,
|
||||
dispatcher: disp,
|
||||
logger: logger.WithName("watcher"),
|
||||
}
|
||||
|
||||
// Регистрируем reconciler для Trigger
|
||||
if err := (&triggerEventReconciler{
|
||||
Client: mgr.GetClient(),
|
||||
dispatcher: disp,
|
||||
logger: w.logger,
|
||||
}).SetupWithManager(mgr); err != nil {
|
||||
return nil, fmt.Errorf("setup reconciler: %w", err)
|
||||
}
|
||||
|
||||
return w, nil
|
||||
}
|
||||
|
||||
// Run запускает manager (блокирует до ctx.Done()).
|
||||
func (w *Watcher) Run(ctx context.Context) error {
|
||||
return w.manager.Start(ctx)
|
||||
}
|
||||
|
||||
// triggerEventReconciler — controller-runtime reconciler только для Trigger{type:event}.
|
||||
type triggerEventReconciler struct {
|
||||
client.Client
|
||||
dispatcher *Dispatcher
|
||||
logger logr.Logger
|
||||
}
|
||||
|
||||
func (r *triggerEventReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
|
||||
tr := &slessv1alpha1.Trigger{}
|
||||
if err := r.Get(ctx, req.NamespacedName, tr); err != nil {
|
||||
if client.IgnoreNotFound(err) == nil {
|
||||
// Триггер удалён — убираем consumer
|
||||
r.dispatcher.Unsubscribe(req.String())
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
// Нас интересуют только event-триггеры
|
||||
if tr.Spec.Type != slessv1alpha1.TriggerTypeEvent {
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
// При удалении — убираем consumer
|
||||
if !tr.DeletionTimestamp.IsZero() {
|
||||
r.dispatcher.Unsubscribe(req.String())
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
// Триггер отключён — не подписываемся
|
||||
if !tr.Spec.Enabled {
|
||||
r.dispatcher.Unsubscribe(req.String())
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
if tr.Spec.Queue == "" {
|
||||
r.logger.Info("event trigger has no queue, skipping", "trigger", req.String())
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
// Строим внутренний URL Service функции.
|
||||
// Оператор при reconcileEvent создаёт Service с именем functionRef в namespace триггера.
|
||||
targetURL := fmt.Sprintf("http://%s.%s.svc.cluster.local:8080/",
|
||||
tr.Spec.FunctionRef, tr.Namespace)
|
||||
|
||||
r.dispatcher.Subscribe(req.String(), tr.Spec.Queue, targetURL)
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
func (r *triggerEventReconciler) SetupWithManager(mgr ctrl.Manager) error {
|
||||
return ctrl.NewControllerManagedBy(mgr).
|
||||
For(&slessv1alpha1.Trigger{}).
|
||||
Complete(r)
|
||||
}
|
||||
Reference in New Issue
Block a user