943 lines
35 KiB
Go
943 lines
35 KiB
Go
// Изменено: 2026-03-23 (deps поддержка, Building poll, rename, kind=service по умолч.)
|
||
// 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 триггера
|
||
// GET /funcs/<namespace>/fn-status/<fn> — прокси: текущий phase+url функции (для UI polling)
|
||
// POST /funcs/<namespace>/create-function — создать функцию через UI (kind=service по умолч.)
|
||
// POST /funcs/<namespace>/save-function/<fn> — обновить код через UI (zip с code+deps)
|
||
// POST /funcs/<namespace>/rename-function/<fn> — переименовать: delete old + create new с тем же кодом
|
||
// 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 "fn-status":
|
||
// GET /funcs/{ns}/fn-status/{fn} — текущий phase+url для UI polling
|
||
if r.Method != http.MethodGet {
|
||
http.Error(w, "method not allowed\n", http.StatusMethodNotAllowed)
|
||
return
|
||
}
|
||
proxyFnStatus(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 "rename-function":
|
||
// POST /funcs/{ns}/rename-function/{fn} — переименовать: delete+create с тем же кодом
|
||
if r.Method != http.MethodPost {
|
||
http.Error(w, "method not allowed\n", http.StatusMethodNotAllowed)
|
||
return
|
||
}
|
||
proxyRenameFunction(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.
|
||
// Создаёт функцию как sless_service (always-on, с URL) если kind не указан.
|
||
// Принимает code (основной файл) + deps (requirements.txt / package.json).
|
||
func proxyCreateFunction(w http.ResponseWriter, r *http.Request, operatorURL, serviceToken, ns string) {
|
||
var req struct {
|
||
Name string `json:"name"`
|
||
Runtime string `json:"runtime"`
|
||
Kind string `json:"kind"` // "service" | "function"; default "service"
|
||
Code string `json:"code"`
|
||
Filename string `json:"filename"`
|
||
Deps string `json:"deps"` // содержимое requirements.txt / package.json
|
||
DepsFile string `json:"deps_filename"` // "requirements.txt" или "package.json"
|
||
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.Kind == "" {
|
||
req.Kind = "service" // always-on по умолчанию — сразу получает URL
|
||
}
|
||
if req.MemoryMB <= 0 {
|
||
req.MemoryMB = 128
|
||
}
|
||
if req.TimeoutSec <= 0 {
|
||
req.TimeoutSec = 30
|
||
}
|
||
|
||
resourceType := "functions"
|
||
if req.Kind == "service" {
|
||
resourceType = "services"
|
||
}
|
||
|
||
// Шаг 1: создаём CRD через оператор (/functions или /services)
|
||
createBody, _ := json.Marshal(map[string]any{
|
||
"name": req.Name,
|
||
"runtime": req.Runtime,
|
||
"entrypoint": req.Entrypoint,
|
||
"memory_mb": req.MemoryMB,
|
||
"timeout_sec": req.TimeoutSec,
|
||
})
|
||
createURL := operatorURL + "/v1/namespaces/" + ns + "/" + resourceType
|
||
createResp, err := operatorRequest(r.Context(), http.MethodPost, createURL, serviceToken, "application/json", bytes.NewReader(createBody))
|
||
if err != nil || (createResp.StatusCode != http.StatusCreated && createResp.StatusCode != http.StatusOK) {
|
||
code := http.StatusBadGateway
|
||
msg := "create 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: упаковываем code + deps в zip и загружаем
|
||
zipFiles := map[string]string{req.Filename: req.Code}
|
||
if req.Deps != "" && req.DepsFile != "" {
|
||
zipFiles[req.DepsFile] = req.Deps
|
||
}
|
||
zipData, err := buildZipFiles(zipFiles)
|
||
if err != nil {
|
||
http.Error(w, `{"error":"zip build failed"}`, http.StatusInternalServerError)
|
||
return
|
||
}
|
||
uploadURL := operatorURL + "/v1/namespaces/" + ns + "/" + resourceType + "/" + req.Name + "/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")
|
||
w.WriteHeader(http.StatusCreated)
|
||
fmt.Fprintf(w, `{"status":"ok","name":%q,"kind":%q}`, req.Name, req.Kind)
|
||
}
|
||
|
||
// proxySaveFunction обрабатывает POST /funcs/{ns}/save-function/{fnName}.
|
||
// Принимает code + deps (опционально), упаковывает в 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"
|
||
Deps string `json:"deps"` // содержимое requirements.txt / package.json
|
||
DepsFile string `json:"deps_filename"` // имя файла зависимостей
|
||
}
|
||
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
|
||
}
|
||
|
||
zipFiles := map[string]string{req.Filename: req.Code}
|
||
if req.Deps != "" && req.DepsFile != "" {
|
||
zipFiles[req.DepsFile] = req.Deps
|
||
}
|
||
zipData, err := buildZipFiles(zipFiles)
|
||
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)
|
||
}
|
||
|
||
// buildZipFiles упаковывает набор текстовых файлов {filename→content} в zip-архив.
|
||
// Используется при создании/редактировании функций через UI.
|
||
// Поддерживает code + deps (requirements.txt / package.json) в одном архиве.
|
||
func buildZipFiles(files map[string]string) ([]byte, error) {
|
||
var buf bytes.Buffer
|
||
zw := zip.NewWriter(&buf)
|
||
for name, content := range files {
|
||
f, err := zw.Create(name)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if _, err := io.WriteString(f, content); err != nil {
|
||
return nil, err
|
||
}
|
||
}
|
||
if err := zw.Close(); err != nil {
|
||
return nil, err
|
||
}
|
||
return buf.Bytes(), nil
|
||
}
|
||
|
||
// proxyFnStatus отдаёт текущий phase+url функции/сервиса — используется для UI polling.
|
||
// Клиент опрашивает каждые 5с пока phase=Building, потом останавливается.
|
||
func proxyFnStatus(w http.ResponseWriter, r *http.Request, operatorURL, serviceToken, ns, fnName string) {
|
||
if !isValidK8sName(fnName) {
|
||
http.Error(w, `{"error":"invalid name"}`, http.StatusBadRequest)
|
||
return
|
||
}
|
||
resourceType := "functions"
|
||
if r.URL.Query().Get("kind") == "service" {
|
||
resourceType = "services"
|
||
}
|
||
resp, err := operatorRequest(r.Context(), http.MethodGet,
|
||
operatorURL+"/v1/namespaces/"+ns+"/"+resourceType+"/"+fnName, serviceToken, "", nil)
|
||
if err != nil {
|
||
http.Error(w, err.Error(), http.StatusBadGateway)
|
||
return
|
||
}
|
||
defer resp.Body.Close()
|
||
b, _ := io.ReadAll(resp.Body)
|
||
w.Header().Set("Content-Type", "application/json")
|
||
w.WriteHeader(resp.StatusCode)
|
||
w.Write(b) //nolint:errcheck
|
||
}
|
||
|
||
// proxyRenameFunction реализует переименование функции: delete old + create new + upload old code.
|
||
// Это атомарная операция с точки зрения UI, но не транзакционная на уровне k8s.
|
||
// При сбое на шагах 3-5 старая функция может остаться (UI покажет ошибку).
|
||
func proxyRenameFunction(w http.ResponseWriter, r *http.Request, operatorURL, serviceToken, ns, oldName string) {
|
||
if !isValidK8sName(oldName) {
|
||
http.Error(w, `{"error":"invalid old name"}`, http.StatusBadRequest)
|
||
return
|
||
}
|
||
var req struct {
|
||
NewName string `json:"new_name"`
|
||
Kind string `json:"kind"`
|
||
}
|
||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil || req.NewName == "" {
|
||
http.Error(w, `{"error":"new_name required"}`, http.StatusBadRequest)
|
||
return
|
||
}
|
||
if !isValidK8sName(req.NewName) {
|
||
http.Error(w, `{"error":"invalid new_name"}`, http.StatusBadRequest)
|
||
return
|
||
}
|
||
resourceType := "functions"
|
||
if req.Kind == "service" {
|
||
resourceType = "services"
|
||
}
|
||
ctx := r.Context()
|
||
baseURL := operatorURL + "/v1/namespaces/" + ns + "/" + resourceType
|
||
|
||
// Шаг 1: получаем конфиг старой функции (runtime, entrypoint, memory, timeout)
|
||
cfgResp, err := operatorRequest(ctx, http.MethodGet, baseURL+"/"+oldName, serviceToken, "", nil)
|
||
if err != nil || cfgResp.StatusCode != http.StatusOK {
|
||
http.Error(w, "get old config failed", http.StatusBadGateway)
|
||
return
|
||
}
|
||
var cfg struct {
|
||
Runtime string `json:"runtime"`
|
||
Entrypoint string `json:"entrypoint"`
|
||
MemoryMB int `json:"memory_mb"`
|
||
TimeoutSec int `json:"timeout_sec"`
|
||
Env map[string]string `json:"env"`
|
||
}
|
||
json.NewDecoder(cfgResp.Body).Decode(&cfg) //nolint:errcheck
|
||
cfgResp.Body.Close()
|
||
|
||
// Шаг 2: получаем исходный код старой функции
|
||
srcResp, _ := operatorRequest(ctx, http.MethodGet, baseURL+"/"+oldName+"/source", serviceToken, "", nil)
|
||
zipFiles := map[string]string{}
|
||
if srcResp != nil && srcResp.StatusCode == http.StatusOK {
|
||
var files []struct {
|
||
Name string `json:"name"`
|
||
Content string `json:"content"`
|
||
Binary bool `json:"binary"`
|
||
}
|
||
if json.NewDecoder(srcResp.Body).Decode(&files) == nil {
|
||
for _, f := range files {
|
||
if !f.Binary {
|
||
zipFiles[f.Name] = f.Content
|
||
}
|
||
}
|
||
}
|
||
srcResp.Body.Close()
|
||
}
|
||
|
||
// Шаг 3: создаём новую функцию с новым именем
|
||
createBody, _ := json.Marshal(map[string]any{
|
||
"name": req.NewName, "runtime": cfg.Runtime, "entrypoint": cfg.Entrypoint,
|
||
"memory_mb": cfg.MemoryMB, "timeout_sec": cfg.TimeoutSec, "env": cfg.Env,
|
||
})
|
||
createResp, err := operatorRequest(ctx, http.MethodPost, baseURL, serviceToken, "application/json", bytes.NewReader(createBody))
|
||
if err != nil || (createResp.StatusCode != http.StatusCreated && createResp.StatusCode != http.StatusOK) {
|
||
msg := "create new failed"
|
||
if createResp != nil {
|
||
b, _ := io.ReadAll(createResp.Body)
|
||
createResp.Body.Close()
|
||
msg = string(b)
|
||
}
|
||
http.Error(w, msg, http.StatusBadGateway)
|
||
return
|
||
}
|
||
createResp.Body.Close()
|
||
|
||
// Шаг 4: загружаем старый код в новую функцию (если есть)
|
||
if len(zipFiles) > 0 {
|
||
if zipData, err := buildZipFiles(zipFiles); err == nil {
|
||
_ = uploadZipToOperator(ctx, baseURL+"/"+req.NewName+"/upload", serviceToken, req.NewName, zipData)
|
||
}
|
||
}
|
||
|
||
// Шаг 5: удаляем старую функцию
|
||
delResp, _ := operatorRequest(ctx, http.MethodDelete, baseURL+"/"+oldName, serviceToken, "", nil)
|
||
if delResp != nil {
|
||
delResp.Body.Close()
|
||
}
|
||
|
||
w.Header().Set("Content-Type", "application/json")
|
||
fmt.Fprintf(w, `{"status":"ok","old":%q,"new":%q,"kind":%q}`, oldName, req.NewName, req.Kind)
|
||
}
|
||
|
||
// 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)
|
||
}
|