Files
sless/services/funcs/main.go
T
Naeel b7fa8acf76 feat(web-console): create/edit/delete functions in UI, v0.1.4
- services/funcs: новые маршруты POST create-function, POST save-function,
  DELETE delete-function — создание/редактирование/удаление через браузер
- index.html: модалка создания с шаблонами Python/Node.js hello world,
  кнопки edit/delete на каждой карточке функции
- sless-funcs-service:v0.1.4 задеплоен
- examples/POSTGRES: удалены логи, .bak, backup tfstate, лишние функции
  оставлены 2 Python (pg-stats, pg-counter) + 2 Node.js (pg-info, js-idempotent)
2026-03-23 10:19:42 +03:00

780 lines
29 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Изменено: 2026-03-23 (добавлены эндпоинты создания/редактирования/удаления функций через UI)
// main.go — глобальный HTTP сервис листинга функций пользователя.
// Развёрнут ОДИН РАЗ в namespace sless; работает для ВСЕХ пользователей.
// Не связан с terraform — деплоится манифестом deployments/k8s/funcs-service.yaml.
//
// Маршруты:
// GET /funcs — usage hint (без токена)
// GET /funcs?token=<jwt> — листинг через JWT токен
// GET /funcs/<namespace> — листинг (Accept: text/html → HTML, иначе plain text)
// GET /funcs/<namespace>/source/<fn> — прокси к оператору: файлы исходного кода (JSON)
// PATCH /funcs/<namespace>/triggers/<n> — прокси к оператору: enable/disable триггера
// POST /funcs/<namespace>/create-function — создать функцию через UI (zip собирается на сервере)
// POST /funcs/<namespace>/save-function/<fn> — сохранить новый код функции через UI
// DELETE /funcs/<namespace>/delete-function/<fn> — удалить функцию через UI
// GET /health — liveness/readiness probe
//
// Env vars:
// SLESS_OPERATOR_URL — URL оператора внутри кластера (default: http://sless-operator.sless.svc.cluster.local:9090)
// SLESS_EXTERNAL_URL — публичный базовый URL для корректных ссылок на функции
// SLESS_EXCLUDE — comma-separated список имён функций, скрытых из листинга
// SLESS_SERVICE_TOKEN — токен сервиса для запросов к оператору (для /funcs/<namespace> без токена юзера)
// PORT — порт сервера (default: 8090)
package main
import (
"archive/zip"
"bytes"
"context"
"crypto/sha256"
_ "embed"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"log"
"mime/multipart"
"net/http"
"os"
"regexp"
"sort"
"strings"
)
//go:embed index.html
var htmlPageTemplate string
// fnResponse — унифицированная запись функции для листинга.
// Kind="function" — job-style (запускается per-request).
// Kind="service" — always-on Deployment (всегда готов к запросам, имеет прямой URL).
type fnResponse struct {
Name string `json:"name"`
Runtime string `json:"runtime"`
Phase string `json:"phase"`
Message string `json:"message"`
CreatedAt string `json:"created_at"`
LastBuiltAt string `json:"last_built_at"`
Kind string `json:"kind"` // "function" | "service"
URL string `json:"url,omitempty"` // прямой URL для sless_service
}
// svcResponse — ответ /v1/namespaces/{ns}/services (подмножество полей оператора).
type svcResponse struct {
Name string `json:"name"`
Runtime string `json:"runtime"`
Phase string `json:"phase"`
Message string `json:"message"`
CreatedAt string `json:"created_at"`
LastBuiltAt string `json:"last_built_at"`
URL string `json:"url"`
}
// trResponse — ответ /v1/namespaces/{ns}/triggers
type trResponse struct {
Name string `json:"name"`
Type string `json:"type"`
FunctionRef string `json:"function"`
Schedule string `json:"schedule"`
Enabled bool `json:"enabled"`
Active bool `json:"active"`
URL string `json:"url"`
}
// pageData — структура данных для HTML-шаблона web-консоли.
type pageData struct {
Namespace string `json:"namespace"`
ExternalURL string `json:"externalURL"`
Functions []fnResponse `json:"functions"`
TriggersByFn map[string][]trResponse `json:"triggersByFn"`
}
// nsRegex — базовая валидация kubernetes namespace (ловеркейс + hyphens).
// Защита от path traversal при подстановке в URL запросов к оператору.
var nsRegex = regexp.MustCompile(`^[a-z0-9]([a-z0-9\-]{0,61}[a-z0-9])?$`)
func main() {
operatorURL := strings.TrimRight(env("SLESS_OPERATOR_URL", "http://sless-operator.sless.svc.cluster.local:9090"), "/")
externalURL := strings.TrimRight(env("SLESS_EXTERNAL_URL", ""), "/")
port := env("PORT", "8090")
serviceToken := os.Getenv("SLESS_SERVICE_TOKEN")
exclude := map[string]bool{}
for _, n := range strings.Split(os.Getenv("SLESS_EXCLUDE"), ",") {
if n = strings.TrimSpace(n); n != "" {
exclude[n] = true
}
}
// /funcs (точное совпадение) — namespace из JWT токена (?token= или Authorization header)
http.HandleFunc("/funcs", handleFuncsToken(operatorURL, externalURL, exclude))
// /funcs/ — namespace в URL пути + sub-paths (source/triggers proxy)
http.HandleFunc("/funcs/", handleFuncsNS(operatorURL, externalURL, serviceToken, exclude))
// /health — для liveness/readiness probe (без auth)
http.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
fmt.Fprintln(w, "ok")
})
log.Printf("sless-funcs-service listening on :%s (operator: %s)", port, operatorURL)
log.Fatal(http.ListenAndServe(":"+port, nil))
}
// handleFuncsToken обрабатывает GET /funcs с JWT токеном (?token= или Authorization header).
func handleFuncsToken(operatorURL, externalURL string, exclude map[string]bool) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
token := r.URL.Query().Get("token")
if token == "" {
token = strings.TrimPrefix(r.Header.Get("Authorization"), "Bearer ")
}
if token == "" {
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
w.WriteHeader(http.StatusUnauthorized)
fmt.Fprintln(w, "Укажите namespace в URL или передайте токен:")
fmt.Fprintln(w, "")
fmt.Fprintf(w, " %s/funcs/<namespace>\n", externalURL)
fmt.Fprintf(w, " %s/funcs?token=<jwt>\n", externalURL)
return
}
sub, err := subFromJWT(token)
if err != nil {
http.Error(w, fmt.Sprintf("invalid token: %s\n", err), http.StatusUnauthorized)
return
}
fetchAndRender(w, r, operatorURL, externalURL, "Bearer "+token, namespaceFromSub(sub), exclude)
}
}
// handleFuncsNS обрабатывает все запросы /funcs/{ns}[/sub/path].
//
// Маршруты:
//
// GET /funcs/{ns} — листинг (HTML если Accept: text/html, иначе plain text)
// GET /funcs/{ns}/source/{fn} — прокси: файлы исходного кода функции
// PATCH /funcs/{ns}/triggers/{name} — прокси: enable/disable триггер
func handleFuncsNS(operatorURL, externalURL, serviceToken string, exclude map[string]bool) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
rest := strings.TrimPrefix(r.URL.Path, "/funcs/")
rest = strings.Trim(rest, "/")
if rest == "" {
http.Error(w, "usage: /funcs/<namespace>\n", http.StatusBadRequest)
return
}
// SplitN до 3 частей: [ns], [ns, subtype] или [ns, subtype, resourceName]
parts := strings.SplitN(rest, "/", 3)
ns := parts[0]
if !nsRegex.MatchString(ns) {
http.Error(w, "invalid namespace\n", http.StatusBadRequest)
return
}
if len(parts) == 3 {
switch parts[1] {
case "source":
if r.Method != http.MethodGet {
http.Error(w, "method not allowed\n", http.StatusMethodNotAllowed)
return
}
proxySourceGet(w, r, operatorURL, serviceToken, ns, parts[2])
return
case "triggers":
if r.Method != http.MethodPatch {
http.Error(w, "method not allowed\n", http.StatusMethodNotAllowed)
return
}
proxyTriggerPatch(w, r, operatorURL, serviceToken, ns, parts[2])
return
case "save-function":
// POST /funcs/{ns}/save-function/{fnName} — обновить код функции через UI
if r.Method != http.MethodPost {
http.Error(w, "method not allowed\n", http.StatusMethodNotAllowed)
return
}
proxySaveFunction(w, r, operatorURL, serviceToken, ns, parts[2])
return
case "delete-function":
// DELETE /funcs/{ns}/delete-function/{fnName} — удалить функцию через UI
if r.Method != http.MethodDelete {
http.Error(w, "method not allowed\n", http.StatusMethodNotAllowed)
return
}
proxyDeleteFunction(w, r, operatorURL, serviceToken, ns, parts[2])
return
}
}
if len(parts) == 2 && parts[1] == "create-function" {
// POST /funcs/{ns}/create-function — создать новую функцию через UI
if r.Method != http.MethodPost {
http.Error(w, "method not allowed\n", http.StatusMethodNotAllowed)
return
}
proxyCreateFunction(w, r, operatorURL, serviceToken, ns)
return
}
if len(parts) > 1 {
http.NotFound(w, r)
return
}
// GET /funcs/{ns} — листинг с авторизацией через serviceToken
authHeader := ""
if serviceToken != "" {
authHeader = "Bearer " + serviceToken
}
fetchAndRender(w, r, operatorURL, externalURL, authHeader, ns, exclude)
}
}
// fetchAndRender загружает functions + services из оператора и отдаёт ответ.
// Оба типа объединяются в единый список — пользователь видит их как «функции».
func fetchAndRender(w http.ResponseWriter, r *http.Request, operatorURL, externalURL, authHeader, ns string, exclude map[string]bool) {
fns, err := apiGet[[]fnResponse](operatorURL+"/v1/namespaces/"+ns+"/functions", authHeader)
if err != nil {
http.Error(w, fmt.Sprintf("operator error (functions): %s\n", err), http.StatusBadGateway)
return
}
for i := range fns {
fns[i].Kind = "function"
}
// sless_service — always-on Deployment + URL. Объединяем с функциями.
svcs, svcErr := apiGet[[]svcResponse](operatorURL+"/v1/namespaces/"+ns+"/services", authHeader)
if svcErr == nil {
for _, svc := range svcs {
fns = append(fns, fnResponse{
Name: svc.Name,
Runtime: svc.Runtime,
Phase: svc.Phase,
Message: svc.Message,
CreatedAt: svc.CreatedAt,
LastBuiltAt: svc.LastBuiltAt,
Kind: "service",
URL: svc.URL,
})
}
}
trs, err := apiGet[[]trResponse](operatorURL+"/v1/namespaces/"+ns+"/triggers", authHeader)
if err != nil {
http.Error(w, fmt.Sprintf("operator error (triggers): %s\n", err), http.StatusBadGateway)
return
}
if acceptsHTML(r) {
serveHTML(w, ns, externalURL, fns, trs, exclude)
} else {
servePlainText(w, ns, externalURL, fns, trs, exclude)
}
}
// acceptsHTML возвращает true если клиент предпочитает HTML (браузер).
// curl без -H "Accept: text/html" получит plain text — поведение совместимо с v0.1.x.
func acceptsHTML(r *http.Request) bool {
return strings.Contains(r.Header.Get("Accept"), "text/html")
}
// serveHTML рендерит web-консоль.
// Данные встраиваются как JSON в <script type="application/json"> — изолировано от HTML-контекста.
// </script> в данных экранируется как <\/script> (валидный JSON escape).
func serveHTML(w http.ResponseWriter, ns, externalURL string, fns []fnResponse, trs []trResponse, exclude map[string]bool) {
trigIdx := buildTriggerIndex(trs)
filtered := filterAndSort(fns, trigIdx, exclude)
pd := pageData{
Namespace: ns,
ExternalURL: externalURL,
Functions: filtered,
TriggersByFn: trigIdx,
}
raw, _ := json.Marshal(pd)
// Предотвращаем преждевременное закрытие тега <script type="application/json">
safeJSON := strings.ReplaceAll(string(raw), "</", `<\/`)
page := strings.ReplaceAll(htmlPageTemplate, "__PAGE_DATA__", safeJSON)
w.Header().Set("Content-Type", "text/html; charset=utf-8")
fmt.Fprint(w, page)
}
// servePlainText рендерит human-readable листинг для терминала (curl).
func servePlainText(w http.ResponseWriter, ns, externalURL string, fns []fnResponse, trs []trResponse, exclude map[string]bool) {
trigIdx := buildTriggerIndex(trs)
type entry struct {
fn fnResponse
httpT []trResponse
cronT []trResponse
}
var items []entry
for _, fn := range fns {
if exclude[fn.Name] {
continue
}
var httpT, cronT []trResponse
for _, tr := range trigIdx[fn.Name] {
switch tr.Type {
case "http":
httpT = append(httpT, tr)
case "cron":
cronT = append(cronT, tr)
}
}
items = append(items, entry{fn, httpT, cronT})
}
sort.Slice(items, func(i, j int) bool {
ai := isActive(items[i].fn.Name, trigIdx)
aj := isActive(items[j].fn.Name, trigIdx)
if ai != aj {
return ai
}
return items[i].fn.Name < items[j].fn.Name
})
sep := strings.Repeat("─", 52)
var sb strings.Builder
for _, it := range items {
fn := it.fn
sb.WriteString(sep + "\n")
sb.WriteString(" " + buildComment(fn, it.httpT, it.cronT) + "\n")
sb.WriteString(fmt.Sprintf(" name: %s\n", fn.Name))
sb.WriteString(fmt.Sprintf(" runtime: %s\n", fn.Runtime))
sb.WriteString(fmt.Sprintf(" phase: %s\n", fn.Phase))
if isActive(fn.Name, trigIdx) {
sb.WriteString(" active: да\n")
} else {
sb.WriteString(" active: нет\n")
}
// Для sless_service — прямой URL из статуса деплоя
if fn.URL != "" {
sb.WriteString(fmt.Sprintf(" url: %s\n", fn.URL))
} else if len(it.httpT) > 0 {
url := it.httpT[0].URL
if externalURL != "" {
url = externalURL + "/fn/" + ns + "/" + fn.Name
}
sb.WriteString(fmt.Sprintf(" url: %s\n", url))
}
if len(it.cronT) > 0 {
sb.WriteString(fmt.Sprintf(" cron: %s\n", it.cronT[0].Schedule))
}
if fn.CreatedAt != "" {
sb.WriteString(fmt.Sprintf(" created: %s\n", fn.CreatedAt))
}
if fn.LastBuiltAt != "" {
sb.WriteString(fmt.Sprintf(" built: %s\n", fn.LastBuiltAt))
}
if fn.Message != "" {
sb.WriteString(fmt.Sprintf(" message: %s\n", fn.Message))
}
}
sb.WriteString(sep + "\n")
sb.WriteString(fmt.Sprintf(" namespace: %s | total: %d\n", ns, len(items)))
sb.WriteString(sep + "\n")
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
fmt.Fprint(w, sb.String())
}
// proxySourceGet проксирует GET /funcs/{ns}/source/{fn}?kind={function|service} → оператор.
// Авторизация — сервисный токен (хранится на сервере, не открывается клиенту).
// kind=service → /v1/.../services/{fn}/source; иначе → /v1/.../functions/{fn}/source
func proxySourceGet(w http.ResponseWriter, r *http.Request, operatorURL, serviceToken, ns, funcName string) {
if !isValidK8sName(funcName) {
http.Error(w, "invalid function name\n", http.StatusBadRequest)
return
}
resourceType := "functions"
if r.URL.Query().Get("kind") == "service" {
resourceType = "services"
}
target := operatorURL + "/v1/namespaces/" + ns + "/" + resourceType + "/" + funcName + "/source"
req, err := http.NewRequestWithContext(r.Context(), http.MethodGet, target, nil)
if err != nil {
http.Error(w, "internal error\n", http.StatusInternalServerError)
return
}
if serviceToken != "" {
req.Header.Set("Authorization", "Bearer "+serviceToken)
}
resp, err := http.DefaultClient.Do(req)
if err != nil {
http.Error(w, fmt.Sprintf("operator error: %s\n", err), http.StatusBadGateway)
return
}
defer resp.Body.Close()
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(resp.StatusCode)
io.Copy(w, resp.Body) //nolint:errcheck
}
// proxyTriggerPatch проксирует PATCH /funcs/{ns}/triggers/{name} → оператор.
// Принимает только поле "enabled" — защита от произвольных изменений с сервисным токеном.
func proxyTriggerPatch(w http.ResponseWriter, r *http.Request, operatorURL, serviceToken, ns, triggerName string) {
if !isValidK8sName(triggerName) {
http.Error(w, "invalid trigger name\n", http.StatusBadRequest)
return
}
// Читаем только поле enabled — остальные поля отбрасываем.
var body struct {
Enabled *bool `json:"enabled"`
}
if err := json.NewDecoder(r.Body).Decode(&body); err != nil || body.Enabled == nil {
http.Error(w, `{"error":"required: {\"enabled\": true/false}"}`, http.StatusBadRequest)
return
}
reqBody, _ := json.Marshal(map[string]bool{"enabled": *body.Enabled})
target := operatorURL + "/v1/namespaces/" + ns + "/triggers/" + triggerName
opReq, err := http.NewRequestWithContext(r.Context(), http.MethodPatch, target, bytes.NewReader(reqBody))
if err != nil {
http.Error(w, "internal error\n", http.StatusInternalServerError)
return
}
opReq.Header.Set("Content-Type", "application/json")
if serviceToken != "" {
opReq.Header.Set("Authorization", "Bearer "+serviceToken)
}
resp, err := http.DefaultClient.Do(opReq)
if err != nil {
http.Error(w, fmt.Sprintf("operator error: %s\n", err), http.StatusBadGateway)
return
}
defer resp.Body.Close()
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(resp.StatusCode)
io.Copy(w, resp.Body) //nolint:errcheck
}
// buildTriggerIndex строит индекс триггеров по имени функции (FunctionRef).
func buildTriggerIndex(trs []trResponse) map[string][]trResponse {
idx := map[string][]trResponse{}
for _, tr := range trs {
idx[tr.FunctionRef] = append(idx[tr.FunctionRef], tr)
}
return idx
}
// filterAndSort фильтрует функции (exclude) и сортирует: активные вверх, затем по имени.
func filterAndSort(fns []fnResponse, trigIdx map[string][]trResponse, exclude map[string]bool) []fnResponse {
result := make([]fnResponse, 0, len(fns))
for _, fn := range fns {
if !exclude[fn.Name] {
result = append(result, fn)
}
}
sort.Slice(result, func(i, j int) bool {
ai := isActive(result[i].Name, trigIdx)
aj := isActive(result[j].Name, trigIdx)
if ai != aj {
return ai
}
return result[i].Name < result[j].Name
})
return result
}
// isActive возвращает true если хотя бы один триггер функции включён и активен.
func isActive(funcName string, trigIdx map[string][]trResponse) bool {
for _, tr := range trigIdx[funcName] {
if tr.Enabled && tr.Active {
return true
}
}
return false
}
// isValidK8sName базовая проверка имени k8s-ресурса: нет path traversal символов.
func isValidK8sName(s string) bool {
return s != "" && !strings.Contains(s, "/") && !strings.Contains(s, "..") && len(s) <= 253
}
func buildComment(fn fnResponse, httpT, cronT []trResponse) string {
if fn.Kind == "service" {
active := "неактивна"
if fn.Phase == "Ready" {
active = "активна"
}
return fmt.Sprintf("always-on сервис (%s) — %s, %s", fn.Runtime, fn.Phase, active)
}
if len(httpT) > 0 {
active := "активна"
if !httpT[0].Active {
active = "неактивна"
}
return fmt.Sprintf("HTTP endpoint (%s) — %s, %s", fn.Runtime, fn.Phase, active)
}
if len(cronT) > 0 {
active := "активна"
if !cronT[0].Active {
active = "неактивна"
}
return fmt.Sprintf("Cron '%s' (%s) — %s, %s", cronT[0].Schedule, fn.Runtime, fn.Phase, active)
}
return fmt.Sprintf("Job/runner без триггера (%s) — %s", fn.Runtime, fn.Phase)
}
// subFromJWT декодирует JWT payload (без проверки подписи) и возвращает sub.
// Подпись не проверяется — trusted perimeter: сервис работает за Ingress.
func subFromJWT(token string) (string, error) {
parts := strings.Split(token, ".")
if len(parts) != 3 {
return "", fmt.Errorf("invalid jwt: expected 3 parts")
}
// base64url без паддинга — добавляем паддинг стандартно
payload := parts[1]
switch len(payload) % 4 {
case 2:
payload += "=="
case 3:
payload += "="
}
data, err := base64.URLEncoding.DecodeString(payload)
if err != nil {
return "", fmt.Errorf("decode payload: %w", err)
}
var claims map[string]any
if err := json.Unmarshal(data, &claims); err != nil {
return "", fmt.Errorf("unmarshal claims: %w", err)
}
sub, ok := claims["sub"].(string)
if !ok || sub == "" {
return "", fmt.Errorf("missing sub claim")
}
return sub, nil
}
// namespaceFromSub — та же логика что в операторе и terraform провайдере.
// SHA256(sub) → первые 8 байт → hex → "sless-{16 hex символов}"
func namespaceFromSub(sub string) string {
hash := sha256.Sum256([]byte(sub))
return fmt.Sprintf("sless-%x", hash[:8])
}
func apiGet[T any](url, authHeader string) (T, error) {
var zero T
req, err := http.NewRequest(http.MethodGet, url, nil)
if err != nil {
return zero, err
}
req.Header.Set("Authorization", authHeader)
resp, err := http.DefaultClient.Do(req)
if err != nil {
return zero, err
}
defer resp.Body.Close()
body, _ := io.ReadAll(resp.Body)
if resp.StatusCode != http.StatusOK {
return zero, fmt.Errorf("status %d: %s", resp.StatusCode, body)
}
if err := json.Unmarshal(body, &zero); err != nil {
return zero, fmt.Errorf("unmarshal: %w", err)
}
return zero, nil
}
func env(key, fallback string) string {
if v := os.Getenv(key); v != "" {
return v
}
return fallback
}
// proxyCreateFunction обрабатывает POST /funcs/{ns}/create-function.
// Принимает JSON с кодом функции, создаёт CRD через оператор и сразу загружает zip с кодом.
// Это позволяет создать функцию целиком за один запрос из UI без использования terraform.
func proxyCreateFunction(w http.ResponseWriter, r *http.Request, operatorURL, serviceToken, ns string) {
var req struct {
Name string `json:"name"`
Runtime string `json:"runtime"`
Code string `json:"code"`
Filename string `json:"filename"`
Entrypoint string `json:"entrypoint"`
MemoryMB int `json:"memory_mb"`
TimeoutSec int `json:"timeout_sec"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, `{"error":"invalid JSON"}`, http.StatusBadRequest)
return
}
if req.Name == "" || req.Runtime == "" || req.Code == "" || req.Filename == "" || req.Entrypoint == "" {
http.Error(w, `{"error":"name, runtime, code, filename, entrypoint are required"}`, http.StatusBadRequest)
return
}
if !isValidK8sName(req.Name) {
http.Error(w, `{"error":"invalid function name"}`, http.StatusBadRequest)
return
}
if req.MemoryMB <= 0 {
req.MemoryMB = 128
}
if req.TimeoutSec <= 0 {
req.TimeoutSec = 30
}
// Шаг 1: создаём Function CRD через оператор
createBody, _ := json.Marshal(map[string]any{
"name": req.Name,
"runtime": req.Runtime,
"entrypoint": req.Entrypoint,
"memory_mb": req.MemoryMB,
"timeout_sec": req.TimeoutSec,
})
createResp, err := operatorRequest(r.Context(), http.MethodPost,
operatorURL+"/v1/namespaces/"+ns+"/functions", serviceToken, "application/json", bytes.NewReader(createBody))
if err != nil || (createResp.StatusCode != http.StatusCreated && createResp.StatusCode != http.StatusOK) {
code := http.StatusBadGateway
msg := "create function CRD failed"
if createResp != nil {
b, _ := io.ReadAll(createResp.Body)
createResp.Body.Close()
msg = string(b)
code = createResp.StatusCode
}
http.Error(w, msg, code)
return
}
createResp.Body.Close()
// Шаг 2: упаковываем код в zip и загружаем в оператор
zipData, err := buildSingleFileZip(req.Filename, req.Code)
if err != nil {
http.Error(w, `{"error":"zip build failed"}`, http.StatusInternalServerError)
return
}
if err := uploadZipToOperator(r.Context(), operatorURL+"/v1/namespaces/"+ns+"/functions/"+req.Name+"/upload", serviceToken, req.Filename, zipData); err != nil {
// Функция создана, но загрузка провалилась — возвращаем ошибку, UI должен показать её
http.Error(w, "upload code: "+err.Error(), http.StatusBadGateway)
return
}
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusCreated)
fmt.Fprintf(w, `{"status":"ok","name":%q}`, req.Name)
}
// proxySaveFunction обрабатывает POST /funcs/{ns}/save-function/{fnName}.
// Принимает JSON с новым кодом, упаковывает в zip и загружает через оператор.
// kind=service → загружает в /services/{name}/upload, иначе → /functions/{name}/upload.
func proxySaveFunction(w http.ResponseWriter, r *http.Request, operatorURL, serviceToken, ns, fnName string) {
if !isValidK8sName(fnName) {
http.Error(w, `{"error":"invalid function name"}`, http.StatusBadRequest)
return
}
var req struct {
Code string `json:"code"`
Filename string `json:"filename"`
Kind string `json:"kind"` // "function" | "service"
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil || req.Code == "" || req.Filename == "" {
http.Error(w, `{"error":"code and filename are required"}`, http.StatusBadRequest)
return
}
zipData, err := buildSingleFileZip(req.Filename, req.Code)
if err != nil {
http.Error(w, `{"error":"zip build failed"}`, http.StatusInternalServerError)
return
}
resourceType := "functions"
if req.Kind == "service" {
resourceType = "services"
}
uploadURL := operatorURL + "/v1/namespaces/" + ns + "/" + resourceType + "/" + fnName + "/upload"
if err := uploadZipToOperator(r.Context(), uploadURL, serviceToken, req.Filename, zipData); err != nil {
http.Error(w, "upload code: "+err.Error(), http.StatusBadGateway)
return
}
w.Header().Set("Content-Type", "application/json")
fmt.Fprintf(w, `{"status":"ok","name":%q}`, fnName)
}
// proxyDeleteFunction обрабатывает DELETE /funcs/{ns}/delete-function/{fnName}.
// Проксирует DELETE к оператору. kind=service → /services/{name}, иначе → /functions/{name}.
func proxyDeleteFunction(w http.ResponseWriter, r *http.Request, operatorURL, serviceToken, ns, fnName string) {
if !isValidK8sName(fnName) {
http.Error(w, `{"error":"invalid function name"}`, http.StatusBadRequest)
return
}
resourceType := "functions"
if r.URL.Query().Get("kind") == "service" {
resourceType = "services"
}
target := operatorURL + "/v1/namespaces/" + ns + "/" + resourceType + "/" + fnName
resp, err := operatorRequest(r.Context(), http.MethodDelete, target, serviceToken, "", nil)
if err != nil {
http.Error(w, "operator error: "+err.Error(), http.StatusBadGateway)
return
}
defer resp.Body.Close()
if resp.StatusCode >= 300 {
b, _ := io.ReadAll(resp.Body)
http.Error(w, string(b), resp.StatusCode)
return
}
w.Header().Set("Content-Type", "application/json")
fmt.Fprintf(w, `{"status":"ok","name":%q}`, fnName)
}
// buildSingleFileZip упаковывает один текстовый файл в zip-архив.
// Используется при создании/редактировании функций через UI — код пишется в браузере.
func buildSingleFileZip(filename, code string) ([]byte, error) {
var buf bytes.Buffer
zw := zip.NewWriter(&buf)
f, err := zw.Create(filename)
if err != nil {
return nil, err
}
if _, err := io.WriteString(f, code); err != nil {
return nil, err
}
if err := zw.Close(); err != nil {
return nil, err
}
return buf.Bytes(), nil
}
// uploadZipToOperator отправляет zip-данные как multipart/form-data field "code" на URL оператора.
func uploadZipToOperator(ctx context.Context, url, token, filename string, zipData []byte) error {
var body bytes.Buffer
mw := multipart.NewWriter(&body)
fw, err := mw.CreateFormFile("code", filename+".zip")
if err != nil {
return err
}
if _, err := fw.Write(zipData); err != nil {
return err
}
mw.Close()
resp, err := operatorRequest(ctx, http.MethodPost, url, token, mw.FormDataContentType(), &body)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode >= 300 {
b, _ := io.ReadAll(resp.Body)
return fmt.Errorf("status %d: %s", resp.StatusCode, b)
}
return nil
}
// operatorRequest выполняет HTTP-запрос к оператору с авторизацией через serviceToken.
func operatorRequest(ctx context.Context, method, url, token, contentType string, body io.Reader) (*http.Response, error) {
req, err := http.NewRequestWithContext(ctx, method, url, body)
if err != nil {
return nil, err
}
if token != "" {
req.Header.Set("Authorization", "Bearer "+token)
}
if contentType != "" {
req.Header.Set("Content-Type", contentType)
}
return http.DefaultClient.Do(req)
}