Files
sless/services/funcs/main.go
T

552 lines
20 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-21 (добавлен вывод sless_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 /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 (
"bytes"
"crypto/sha256"
_ "embed"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"log"
"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
}
}
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
}