Files
fission-console/console/main.go
T
Naeel 01df1498d6 feat: big-suite terraform example + provider nodejs zip fix
- examples/big-suite: E-Commerce 10 functions, 6 layers, real depends_on
- provider: nodejs packages now wrapped in ESM zip (buildJSDeployZip)
- provider: loadPackageLiteral reads raw file, nodejs zip via source_dir
- all 10 functions verified working: python/node/ruby/go runtimes
2026-04-20 11:38:45 +03:00

1772 lines
60 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.
package main
import (
"archive/zip"
"bytes"
"context"
"crypto/sha256"
"encoding/base64"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"log"
"net"
"net/http"
"os"
"regexp"
"sort"
"strconv"
"strings"
"sync"
"time"
"unicode/utf8"
"fission-console/ui"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/dynamic"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/clientcmd"
)
var (
environmentGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "environments"}
packageGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "packages"}
functionGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "functions"}
httpTrigGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "httptriggers"}
timeTrigGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "timetriggers"}
namespaceGVR = schema.GroupVersionResource{Group: "", Version: "v1", Resource: "namespaces"}
deploymentGVR = schema.GroupVersionResource{Group: "apps", Version: "v1", Resource: "deployments"}
)
const defaultSATokenPath = "/var/run/secrets/kubernetes.io/serviceaccount/token"
var deckAPIs = map[string]string{
"prod": "https://deck-api.ngcloud.ru/api/v1",
"dev": "https://deck-api-dev.ngcloud.ru/api/v1",
"test": "https://deck-api-test.ngcloud.ru/api/v1",
}
type server struct {
dyn dynamic.Interface
ns string
routerURL string
http *http.Client
saTokenPath string
invokeTimeout time.Duration
testMode bool // FISSION_TEST_MODE=true — пропускает deck auth, X-Test-Sub задаёт user
authUser string
authPass string
tokenMu sync.Mutex
cachedJWT string
tokenExpAt time.Time
tokenCache sync.Map
}
type createFunctionRequest struct {
Name string `json:"name"`
Language string `json:"language"`
Environment string `json:"environment"`
Code string `json:"code"`
Entrypoint string `json:"entrypoint"`
Route string `json:"route"`
Methods []string `json:"methods"`
TTL string `json:"ttl"` // e.g. "24h", "7d" — пустое = функция не протухает
}
type langEnvDef struct {
Image string
BuilderImage string
Version int // 0 defaults to 3
}
var langEnvMap = map[string]langEnvDef{
"python": {Image: "ghcr.io/fission/python-env"},
"nodejs": {Image: "ghcr.io/fission/node-env"},
"go": {Image: "ghcr.io/fission/go-env", BuilderImage: "ghcr.io/fission/go-builder"},
"php": {Image: "ghcr.io/fission/php-env"},
"ruby": {Image: "ghcr.io/fission/ruby-env"},
"perl": {Image: "ghcr.io/fission/perl-env", Version: 1},
}
type updateCodeRequest struct {
Code string `json:"code"`
}
func main() {
kubeconfig := strings.TrimSpace(os.Getenv("KUBECONFIG"))
namespace := envDefault("FISSION_NAMESPACE", "default")
routerURL := strings.TrimRight(envDefault("FISSION_ROUTER_URL", "http://router.fission.svc.cluster.local"), "/")
port := envDefault("PORT", "8090")
httpTimeout := envDurationDefault("FISSION_HTTP_TIMEOUT", 30*time.Second)
invokeTimeout := envDurationDefault("FISSION_INVOKE_TIMEOUT", 20*time.Second)
cfg, err := buildConfig(kubeconfig)
if err != nil {
log.Fatalf("build kube config: %v", err)
}
dyn, err := dynamic.NewForConfig(cfg)
if err != nil {
log.Fatalf("create dynamic client: %v", err)
}
authUser := envDefault("FISSION_AUTH_USERNAME", "")
authPass := envDefault("FISSION_AUTH_PASSWORD", "")
saTokenPath := envDefault("SA_TOKEN_PATH", defaultSATokenPath)
s := &server{
dyn: dyn,
ns: namespace,
routerURL: routerURL,
http: &http.Client{Timeout: httpTimeout},
saTokenPath: saTokenPath,
invokeTimeout: invokeTimeout,
authUser: authUser,
authPass: authPass,
testMode: os.Getenv("FISSION_TEST_MODE") == "true",
}
mux := http.NewServeMux()
mux.HandleFunc("/health", func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
_, _ = w.Write([]byte("ok\n"))
})
mux.HandleFunc("/console/health", func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
_, _ = w.Write([]byte("ok\n"))
})
uiHandler := http.StripPrefix("/console", ui.Handler())
mux.Handle("/console", uiHandler)
mux.Handle("/console/", uiHandler)
mux.HandleFunc("/api/environments", s.handleList(environmentGVR))
mux.HandleFunc("/api/packages", s.handleList(packageGVR))
mux.HandleFunc("/api/functions", s.handleFunctionsRoot)
mux.HandleFunc("/api/functions/", s.handleFunctionsAction)
mux.HandleFunc("/api/httptriggers", s.handleList(httpTrigGVR))
mux.HandleFunc("/api/timetriggers", s.handleList(timeTrigGVR))
auth := func(h http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
var ns string
if s.testMode {
// TEST_MODE: deck API не вызывается.
// X-Test-Sub задаёт произвольный sub → разные namespace-ы для тестирования.
sub := strings.TrimSpace(r.Header.Get("X-Test-Sub"))
if sub == "" {
writeJSONError(w, http.StatusUnauthorized, "test mode: X-Test-Sub required")
return
}
h32 := sha256.Sum256([]byte(sub))
ns = "fission-" + hex.EncodeToString(h32[:8])
} else {
token := strings.TrimSpace(r.Header.Get("X-Auth-Token"))
env := strings.TrimSpace(strings.ToLower(r.Header.Get("X-Auth-Env")))
if _, ok := deckAPIs[env]; !ok {
env = "test"
}
if token == "" {
writeJSONError(w, http.StatusUnauthorized, "unauthorized")
return
}
if err := s.validateDeckToken(token, env); err != nil {
writeJSONError(w, http.StatusUnauthorized, "unauthorized")
return
}
var err error
ns, err = namespaceFromJWT(token)
if err != nil {
log.Printf("namespaceFromJWT: %v", err)
ns = s.ns
}
}
ctx := context.WithValue(r.Context(), ctxKeyNS{}, ns)
h(w, r.WithContext(ctx))
}
}
mux.HandleFunc("/console/api/auth", s.handleAuth)
mux.HandleFunc("/console/api/environments", auth(s.handleList(environmentGVR)))
mux.HandleFunc("/console/api/packages", auth(s.handleList(packageGVR)))
mux.HandleFunc("/console/api/functions", auth(s.handleFunctionsRoot))
mux.HandleFunc("/console/api/functions/", auth(s.handleFunctionsAction))
mux.HandleFunc("/console/api/httptriggers", auth(s.handleList(httpTrigGVR)))
mux.HandleFunc("/console/api/timetriggers", auth(s.handleList(timeTrigGVR)))
httpServer := &http.Server{
Addr: ":" + port,
Handler: withSecurityHeaders(withCORS(logRequests(mux))),
ReadHeaderTimeout: 10 * time.Second,
}
log.Printf("fission-console listening on :%s (namespace=%s)", port, namespace)
s.startExpiryReaper(envDurationDefault("REAPER_INTERVAL", 5*time.Minute))
log.Fatal(httpServer.ListenAndServe())
}
func (s *server) handleFunctionsRoot(w http.ResponseWriter, r *http.Request) {
switch r.Method {
case http.MethodGet:
s.handleList(functionGVR)(w, r)
case http.MethodPost:
s.handleCreateFunction(w, r)
default:
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
}
}
func (s *server) handleFunctionsAction(w http.ResponseWriter, r *http.Request) {
path := strings.TrimPrefix(r.URL.Path, "/api/functions/")
if path == r.URL.Path {
path = strings.TrimPrefix(r.URL.Path, "/console/api/functions/")
}
path = strings.Trim(path, "/")
if path == "" {
http.NotFound(w, r)
return
}
parts := strings.Split(path, "/")
name := strings.TrimSpace(parts[0])
if name == "" {
writeJSONError(w, http.StatusBadRequest, "function name is required")
return
}
if len(parts) == 1 {
switch r.Method {
case http.MethodGet:
s.handleGetFunction(w, r, name)
case http.MethodDelete:
s.handleDeleteFunction(w, r, name)
default:
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
}
return
}
if len(parts) == 2 && parts[1] == "code" && r.Method == http.MethodPut {
s.handleUpdateFunctionCode(w, r, name)
return
}
if len(parts) == 2 && parts[1] == "invoke" && r.Method == http.MethodPost {
s.handleInvokeFunction(w, r, name)
return
}
http.NotFound(w, r)
}
// validFuncName — RFC 1123 subdomain label: lowercase alphanumeric + hyphens, no leading/trailing hyphen, max 63 chars.
var validFuncName = regexp.MustCompile(`^[a-z0-9]([a-z0-9-]*[a-z0-9])?$`)
// maxCodeSize — максимальный размер кода функции (1 MB).
const maxCodeSize = 1 << 20
func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) {
var req createFunctionRequest
ns := s.userNS(r)
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("decode request: %v", err))
return
}
// Гарантируем что namespace + RBAC созданы до любых операций с ресурсами.
// handleAuth делает это при логине, но в test mode или при прямом вызове API
// namespace может отсутствовать — создаём idempotent.
nsCtx, nsCancel := context.WithTimeout(r.Context(), 30*time.Second)
defer nsCancel()
if err := s.ensureUserNamespace(nsCtx, ns); err != nil {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure namespace: %v", err))
return
}
req.Name = strings.TrimSpace(req.Name)
req.Language = strings.TrimSpace(req.Language)
req.Environment = strings.TrimSpace(req.Environment)
req.Code = strings.TrimSpace(req.Code)
req.Entrypoint = strings.TrimSpace(req.Entrypoint)
req.Route = strings.TrimSpace(req.Route)
// Валидация имени: RFC 1123 label, максимум 57 символов.
// Ограничение 57 (не 63): самый длинный суффикс — "-route" (HTTPTrigger) = 6 символов.
// 63 - 6 = 57. Fission webhook требует все объекты <= 63 символов.
if req.Name != "" && (!validFuncName.MatchString(req.Name) || len(req.Name) > 57) {
writeJSONError(w, http.StatusBadRequest, "invalid function name: must match ^[a-z0-9]([a-z0-9-]*[a-z0-9])?$ and be <= 57 chars")
return
}
// Лимит размера кода
if len(req.Code) > maxCodeSize {
writeJSONError(w, http.StatusBadRequest, "code exceeds 1MB limit")
return
}
// Resolve language → environment (lazy creation).
// ensureEnvironment создаёт Environment CRD если не существует — Fission увидит и поднимет pool pod.
if req.Language != "" {
envCtx, envCancel := context.WithTimeout(r.Context(), 15*time.Second)
defer envCancel()
envName, err := s.ensureEnvironment(envCtx, ns, req.Language)
if err != nil {
// unsupported language — это клиентская ошибка → 400
if strings.Contains(err.Error(), "unsupported language") {
writeJSONError(w, http.StatusBadRequest, err.Error())
} else {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure environment: %v", err))
}
return
}
req.Environment = envName
}
if req.Name == "" || req.Environment == "" || req.Code == "" {
writeJSONError(w, http.StatusBadRequest, "name, environment/language and code are required")
return
}
if req.Entrypoint == "" {
req.Entrypoint = defaultEntrypoint(req.Language)
}
if req.Route == "" {
// namespace-prefix route to avoid collisions between users
nsShort := s.userNS(r)
if len(nsShort) > 12 {
nsShort = nsShort[len(nsShort)-12:]
}
req.Route = "/" + nsShort + "/" + req.Name
}
if !strings.HasPrefix(req.Route, "/") {
req.Route = "/" + req.Route
}
req.Methods = normalizeMethods(req.Methods)
ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second)
defer cancel()
if _, err := s.dyn.Resource(environmentGVR).Namespace(s.userNS(r)).Get(ctx, req.Environment, metav1.GetOptions{}); err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("environment %q not found: %v", req.Environment, err))
return
}
pkgName := req.Name + "-pkg"
triggerName := req.Name + "-route"
methodValues := make([]any, 0, len(req.Methods))
for _, method := range req.Methods {
methodValues = append(methodValues, method)
}
// Build the package spec: Go uses source archive (builder), others use literal deployment
var pkgSpec map[string]any
if req.Language == "go" {
srcZip, err := s.buildGoSourceZip(req.Code)
if err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("build go source archive: %v", err))
return
}
literal := base64.StdEncoding.EncodeToString(srcZip)
pkgSpec = map[string]any{
"source": map[string]any{
"type": "literal",
"literal": literal,
},
"deployment": map[string]any{},
"environment": map[string]any{
"name": req.Environment,
"namespace": ns,
},
"buildcommand": "build",
}
} else {
var deployBytes []byte
if req.Language == "nodejs" {
zipBytes, zipErr := buildJSDeployZip(req.Code)
if zipErr != nil {
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build nodejs archive: %v", zipErr))
return
}
deployBytes = zipBytes
} else {
deployBytes = []byte(req.Code)
}
literal := base64.StdEncoding.EncodeToString(deployBytes)
pkgSpec = map[string]any{
"deployment": map[string]any{
"type": "literal",
"literal": literal,
},
"environment": map[string]any{
"name": req.Environment,
"namespace": ns,
},
"source": map[string]any{},
}
}
pkg := &unstructured.Unstructured{Object: map[string]any{
"apiVersion": "fission.io/v1",
"kind": "Package",
"metadata": map[string]any{
"name": pkgName,
"namespace": ns,
},
"spec": pkgSpec,
}}
// Парсим TTL ДО создания k8s ресурсов — чтобы невалидный TTL не оставлял мусор
fnAnnotations := map[string]any{}
fnAnnotations["fission-console/language"] = req.Language
if req.TTL != "" {
expiresAt, ttlErr := parseTTL(req.TTL)
if ttlErr != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("invalid ttl %q: %v", req.TTL, ttlErr))
return
}
fnAnnotations["fission-console/expires-at"] = expiresAt.UTC().Format(time.RFC3339)
}
if _, err := s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Create(ctx, pkg, metav1.CreateOptions{}); err != nil {
if apierrors.IsAlreadyExists(err) {
writeJSONError(w, http.StatusConflict, fmt.Sprintf("function %q already exists", req.Name))
return
}
if apierrors.IsInvalid(err) {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("invalid function spec: %v", err))
return
}
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create package: %v", err))
return
}
fn := &unstructured.Unstructured{Object: map[string]any{
"apiVersion": "fission.io/v1",
"kind": "Function",
"metadata": map[string]any{
"name": req.Name,
"namespace": ns,
"annotations": fnAnnotations,
},
"spec": map[string]any{
"environment": map[string]any{
"name": req.Environment,
"namespace": ns,
},
"InvokeStrategy": map[string]any{
"ExecutionStrategy": map[string]any{"ExecutorType": "poolmgr"},
"StrategyType": "execution",
},
"package": map[string]any{
"packageref": map[string]any{"name": pkgName, "namespace": s.userNS(r)},
"functionName": req.Entrypoint,
},
},
}}
if _, err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).Create(ctx, fn, metav1.CreateOptions{}); err != nil {
_ = s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Delete(ctx, pkgName, metav1.DeleteOptions{})
if apierrors.IsAlreadyExists(err) {
writeJSONError(w, http.StatusConflict, fmt.Sprintf("function %q already exists", req.Name))
return
}
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create function: %v", err))
return
}
httpTrigger := &unstructured.Unstructured{Object: map[string]any{
"apiVersion": "fission.io/v1",
"kind": "HTTPTrigger",
"metadata": map[string]any{
"name": triggerName,
"namespace": ns,
},
"spec": map[string]any{
"relativeurl": req.Route,
"methods": methodValues,
"createingress": true,
"functionref": map[string]any{
"type": "name",
"name": req.Name,
},
},
}}
if _, err := s.dyn.Resource(httpTrigGVR).Namespace(s.userNS(r)).Create(ctx, httpTrigger, metav1.CreateOptions{}); err != nil {
_ = s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).Delete(ctx, req.Name, metav1.DeleteOptions{})
_ = s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Delete(ctx, pkgName, metav1.DeleteOptions{})
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create httptrigger: %v", err))
return
}
writeAnyJSON(w, http.StatusCreated, map[string]any{
"name": req.Name,
"package": pkgName,
"httptrigger": triggerName,
"route": req.Route,
"expires_at": fnAnnotations["fission-console/expires-at"],
})
}
func (s *server) buildLangEnvironment(name string, ns string, def langEnvDef) *unstructured.Unstructured {
envVersion := int64(3)
if def.Version != 0 {
envVersion = int64(def.Version)
}
spec := map[string]any{
"version": envVersion,
"runtime": map[string]any{
"image": def.Image,
},
"poolsize": int64(1),
}
if def.BuilderImage != "" {
spec["builder"] = map[string]any{
"image": def.BuilderImage,
"command": "build",
}
}
return &unstructured.Unstructured{Object: map[string]any{
"apiVersion": "fission.io/v1",
"kind": "Environment",
"metadata": map[string]any{
"name": name,
"namespace": ns,
},
"spec": spec,
}}
}
// ensureEnvironment создаёт Environment CRD для языка lang в namespace ns если не существует.
// Вся Fission-специфичная конструкция объекта изолирована в buildLangEnvironment —
// при обновлении Fission (изменение схемы CRD) меняем только там.
func (s *server) ensureEnvironment(ctx context.Context, ns, lang string) (string, error) {
langDef, ok := langEnvMap[lang]
if !ok {
return "", fmt.Errorf("unsupported language: %q", lang)
}
envName := "console-" + lang + "-env"
_, getErr := s.dyn.Resource(environmentGVR).Namespace(ns).Get(ctx, envName, metav1.GetOptions{})
if getErr == nil {
return envName, nil // уже существует
}
if !apierrors.IsNotFound(getErr) {
return "", fmt.Errorf("check environment %q: %w", envName, getErr)
}
env := s.buildLangEnvironment(envName, ns, langDef)
if _, createErr := s.dyn.Resource(environmentGVR).Namespace(ns).Create(ctx, env, metav1.CreateOptions{}); createErr != nil && !apierrors.IsAlreadyExists(createErr) {
return "", fmt.Errorf("create environment %q: %w", envName, createErr)
}
log.Printf("ensureEnvironment: created %s/%s", ns, envName)
return envName, nil
}
// cleanupEnvironmentIfUnused удаляет Environment CRD если ни одна функция в namespace его не использует.
// Fission увидит удаление → уберёт pool Deployment → поды умирают.
// Не блокирующий: ошибки логируются, не возвращаются вызывающему коду.
func (s *server) cleanupEnvironmentIfUnused(ctx context.Context, ns, envName string) {
functions, err := s.dyn.Resource(functionGVR).Namespace(ns).List(ctx, metav1.ListOptions{})
if err != nil {
log.Printf("cleanupEnvironmentIfUnused: list functions in %s: %v", ns, err)
return
}
for _, fn := range functions.Items {
fnEnv, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name")
if fnEnv == envName {
return // ещё используется хотя бы одной функцией
}
}
// Ни одна функция не ссылается на этот environment — удаляем
if delErr := s.dyn.Resource(environmentGVR).Namespace(ns).Delete(ctx, envName, metav1.DeleteOptions{}); delErr != nil && !apierrors.IsNotFound(delErr) {
log.Printf("cleanupEnvironmentIfUnused: delete env %s/%s: %v", ns, envName, delErr)
return
}
log.Printf("cleanupEnvironmentIfUnused: deleted unused env %s/%s", ns, envName)
}
// parseTTL парсит строку TTL и возвращает время истечения.
// Поддерживаемые форматы: Go duration (1h, 30m, 24h) и дни (1d, 7d, 30d).
// Суффикс "d" не поддерживается стандартным time.ParseDuration — обрабатываем отдельно.
func parseTTL(ttl string) (time.Time, error) {
if strings.HasSuffix(ttl, "d") {
days, err := strconv.Atoi(strings.TrimSuffix(ttl, "d"))
if err != nil || days <= 0 {
return time.Time{}, fmt.Errorf("invalid days value: %q", ttl)
}
return time.Now().Add(time.Duration(days) * 24 * time.Hour), nil
}
d, err := time.ParseDuration(ttl)
if err != nil {
return time.Time{}, err
}
if d <= 0 {
return time.Time{}, fmt.Errorf("ttl must be positive")
}
return time.Now().Add(d), nil
}
// startExpiryReaper запускает фоновый goroutine для удаления функций с истёкшим TTL.
// Интервал задаётся через env REAPER_INTERVAL (default 5m).
// При удалении функции вызывает cleanupEnvironmentIfUnused — поды умирают когда language больше не используется.
func (s *server) startExpiryReaper(interval time.Duration) {
go func() {
ticker := time.NewTicker(interval)
defer ticker.Stop()
log.Printf("expiryReaper: started, interval=%v", interval)
for range ticker.C {
s.runExpiryReap()
}
}()
}
// runExpiryReap обходит все namespace под управлением fission-console и удаляет протухшие функции.
func (s *server) runExpiryReap() {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
defer cancel()
// Находим только наши namespace-ы по метке которую мы ставим при создании
nsList, err := s.dyn.Resource(namespaceGVR).List(ctx, metav1.ListOptions{
LabelSelector: "managed-by=fission-console",
})
if err != nil {
log.Printf("expiryReaper: list namespaces: %v", err)
return
}
now := time.Now().UTC()
for _, ns := range nsList.Items {
s.reapExpiredFunctionsInNS(ctx, ns.GetName(), now)
}
}
// reapExpiredFunctionsInNS удаляет протухшие функции в конкретном namespace.
// Для каждой удалённой функции вызывает cleanupEnvironmentIfUnused.
// Также удаляет orphan packages — пакеты у которых нет соответствующей функции.
func (s *server) reapExpiredFunctionsInNS(ctx context.Context, ns string, now time.Time) {
functions, err := s.dyn.Resource(functionGVR).Namespace(ns).List(ctx, metav1.ListOptions{})
if err != nil {
log.Printf("expiryReaper: list functions in %s: %v", ns, err)
return
}
// Строим множество имён существующих функций для поиска orphan packages
activeFunctions := make(map[string]struct{}, len(functions.Items))
for _, fn := range functions.Items {
activeFunctions[fn.GetName()] = struct{}{}
}
for _, fn := range functions.Items {
expiresAtStr, _, _ := unstructured.NestedString(fn.Object, "metadata", "annotations", "fission-console/expires-at")
if expiresAtStr == "" {
continue // нет TTL — функция живёт вечно
}
expiresAt, parseErr := time.Parse(time.RFC3339, expiresAtStr)
if parseErr != nil {
log.Printf("expiryReaper: parse expires-at for %s/%s: %v", ns, fn.GetName(), parseErr)
continue
}
if now.Before(expiresAt) {
continue // ещё не протухла
}
// Функция протухла — удаляем всё
fnName := fn.GetName()
envName, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name")
pkgName, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name")
log.Printf("expiryReaper: deleting expired function %s/%s (expired %s ago)", ns, fnName, now.Sub(expiresAt).Round(time.Second))
// Удаляем HTTPTrigger-ы ссылающиеся на эту функцию
triggers, tErr := s.dyn.Resource(httpTrigGVR).Namespace(ns).List(ctx, metav1.ListOptions{})
if tErr == nil {
for _, trig := range triggers.Items {
refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name")
if refName == fnName {
_ = s.dyn.Resource(httpTrigGVR).Namespace(ns).Delete(ctx, trig.GetName(), metav1.DeleteOptions{})
}
}
}
_ = s.dyn.Resource(functionGVR).Namespace(ns).Delete(ctx, fnName, metav1.DeleteOptions{})
if pkgName != "" {
_ = s.dyn.Resource(packageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{})
}
// Убираем pool pod если язык больше не используется
if envName != "" {
s.cleanupEnvironmentIfUnused(ctx, ns, envName)
}
}
// Сканируем orphan packages — пакеты без соответствующей функции
// (могут остаться если под упал в середине удаления)
packages, pkgListErr := s.dyn.Resource(packageGVR).Namespace(ns).List(ctx, metav1.ListOptions{})
if pkgListErr == nil {
for _, pkg := range packages.Items {
pkgName := pkg.GetName()
// Конвенция именования: {fn-name}-pkg
if !strings.HasSuffix(pkgName, "-pkg") {
continue
}
fnName := strings.TrimSuffix(pkgName, "-pkg")
if _, exists := activeFunctions[fnName]; !exists {
log.Printf("expiryReaper: deleting orphan package %s/%s (no matching function)", ns, pkgName)
_ = s.dyn.Resource(packageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{})
}
}
}
}
func (s *server) buildGoSourceZip(code string) ([]byte, error) {
var buf bytes.Buffer
zw := zip.NewWriter(&buf)
fw, err := zw.Create("handler.go")
if err != nil {
return nil, err
}
if _, err := fw.Write([]byte(code)); err != nil {
return nil, err
}
goMod := "module github.com/user/fn\n\ngo 1.23\n"
fw2, err := zw.Create("go.mod")
if err != nil {
return nil, err
}
if _, err := fw2.Write([]byte(goMod)); err != nil {
return nil, err
}
if err := zw.Close(); err != nil {
return nil, err
}
return buf.Bytes(), nil
}
func buildJSDeployZip(code string) ([]byte, error) {
var buf bytes.Buffer
zw := zip.NewWriter(&buf)
// package.json: объявляем ESM тип чтобы Node.js трактовал .js как ESM модуль
pkgfw, err := zw.Create("package.json")
if err != nil {
return nil, err
}
if _, err := pkgfw.Write([]byte(`{"type":"module"}`)); err != nil {
return nil, err
}
// main.js — ESM wrapper + код пользователя инлайн через new Function
// new Function безопасно изолирует module/exports от глобального контекста
codeJSON, err := json.Marshal(code)
if err != nil {
return nil, fmt.Errorf("marshal user code: %w", err)
}
wrapper := fmt.Sprintf(`const __mod = { exports: {} };
(new Function('module', 'exports', %s))(__mod, __mod.exports);
const _fn = __mod.exports;
export default async function(ctx) {
const fn = typeof _fn === 'function' ? _fn : (_fn.default || _fn.handler || _fn.main);
if (!fn) throw new Error('no exported function found in user code');
const result = await fn(ctx);
if (!result) return { status: 200, body: '' };
if (typeof result.status !== 'undefined') return result;
return { status: 200, ...result };
}
`, string(codeJSON))
fw, err := zw.Create("main.js")
if err != nil {
return nil, err
}
if _, err := fw.Write([]byte(wrapper)); err != nil {
return nil, err
}
if err := zw.Close(); err != nil {
return nil, err
}
return buf.Bytes(), nil
}
// buildScriptZip wraps code into a zip file with the given filename.
// Used for PHP, Ruby, Perl where the environment requires a named source file.
func buildScriptZip(code, filename string) ([]byte, error) {
var buf bytes.Buffer
zw := zip.NewWriter(&buf)
fw, err := zw.Create(filename)
if err != nil {
return nil, err
}
if _, err := fw.Write([]byte(code)); err != nil {
return nil, err
}
if err := zw.Close(); err != nil {
return nil, err
}
return buf.Bytes(), nil
}
func defaultEntrypoint(lang string) string {
switch lang {
case "nodejs":
return "main"
case "php":
return "main.php::handler"
case "ruby":
return "handler"
case "perl":
return "handler"
case "go":
return "Handler"
default:
return "main.main"
}
}
func (s *server) handleGetFunction(w http.ResponseWriter, r *http.Request, name string) {
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second)
defer cancel()
fn, err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).Get(ctx, name, metav1.GetOptions{})
if err != nil {
status := http.StatusBadGateway
if apierrors.IsNotFound(err) {
status = http.StatusNotFound
}
writeJSONError(w, status, fmt.Sprintf("get function %q: %v", name, err))
return
}
packageName, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name")
environment, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name")
entrypoint, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "functionName")
code := ""
if packageName != "" {
pkg, pkgErr := s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Get(ctx, packageName, metav1.GetOptions{})
if pkgErr == nil {
code = s.extractPackageSourceCode(ctx, pkg)
}
}
route := ""
methods := []string{}
triggers, trigErr := s.dyn.Resource(httpTrigGVR).Namespace(s.userNS(r)).List(ctx, metav1.ListOptions{})
if trigErr == nil {
for _, trig := range triggers.Items {
refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name")
if refName != name {
continue
}
route, _, _ = unstructured.NestedString(trig.Object, "spec", "relativeurl")
methods, _, _ = unstructured.NestedStringSlice(trig.Object, "spec", "methods")
break
}
}
writeAnyJSON(w, http.StatusOK, map[string]any{
"name": name,
"namespace": s.userNS(r),
"environment": environment,
"package": packageName,
"entrypoint": entrypoint,
"code": code,
"route": route,
"methods": methods,
"raw": fn.Object,
})
}
func (s *server) handleUpdateFunctionCode(w http.ResponseWriter, r *http.Request, name string) {
var req updateCodeRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("decode request: %v", err))
return
}
req.Code = strings.TrimSpace(req.Code)
if req.Code == "" {
writeJSONError(w, http.StatusBadRequest, "code is required")
return
}
ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second)
defer cancel()
fn, err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).Get(ctx, name, metav1.GetOptions{})
if err != nil {
status := http.StatusBadGateway
if apierrors.IsNotFound(err) {
status = http.StatusNotFound
}
writeJSONError(w, status, fmt.Sprintf("get function %q: %v", name, err))
return
}
pkgName, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name")
if pkgName == "" {
writeJSONError(w, http.StatusBadGateway, "function has no package reference")
return
}
pkg, err := s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Get(ctx, pkgName, metav1.GetOptions{})
if err != nil {
status := http.StatusBadGateway
if apierrors.IsNotFound(err) {
status = http.StatusNotFound
}
writeJSONError(w, status, fmt.Sprintf("get package %q: %v", pkgName, err))
return
}
lang, _, _ := unstructured.NestedString(fn.Object, "metadata", "annotations", "fission-console/language")
var deployBytes []byte
if lang == "nodejs" {
zipBytes, zipErr := buildJSDeployZip(req.Code)
if zipErr != nil {
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build nodejs archive: %v", zipErr))
return
}
deployBytes = zipBytes
} else {
deployBytes = []byte(req.Code)
}
literal := base64.StdEncoding.EncodeToString(deployBytes)
if err := unstructured.SetNestedField(pkg.Object, literal, "spec", "deployment", "literal"); err != nil {
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set package literal: %v", err))
return
}
if _, err := s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Update(ctx, pkg, metav1.UpdateOptions{}); err != nil {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("update package %q: %v", pkgName, err))
return
}
updatedPkg, err := s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Get(ctx, pkgName, metav1.GetOptions{})
if err != nil {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get updated package %q: %v", pkgName, err))
return
}
if err := unstructured.SetNestedField(fn.Object, updatedPkg.GetResourceVersion(), "spec", "package", "packageref", "resourceversion"); err != nil {
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set function package resourceversion: %v", err))
return
}
if _, err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).Update(ctx, fn, metav1.UpdateOptions{}); err != nil {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("update function %q package ref: %v", name, err))
return
}
writeAnyJSON(w, http.StatusOK, map[string]any{"updated": true, "package": pkgName, "package_resourceversion": updatedPkg.GetResourceVersion()})
}
func (s *server) handleInvokeFunction(w http.ResponseWriter, r *http.Request, name string) {
bodyBytes, err := io.ReadAll(r.Body)
if err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("read request body: %v", err))
return
}
if len(bytes.TrimSpace(bodyBytes)) == 0 {
bodyBytes = []byte("{}")
}
invokeTimeout := s.invokeTimeout
if invokeTimeout <= 0 {
invokeTimeout = 20 * time.Second
}
ctx, cancel := context.WithTimeout(r.Context(), invokeTimeout)
defer cancel()
// check function exists in k8s before invoke
if _, err2 := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).Get(ctx, name, metav1.GetOptions{}); err2 != nil {
if apierrors.IsNotFound(err2) {
writeJSONError(w, http.StatusNotFound, fmt.Sprintf("function %q not found", name))
return
}
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get function %q: %v", name, err2))
return
}
invokeURL := fmt.Sprintf("%s/fission-function/v2/functions/%s", s.routerURL, name)
invokeMethod := http.MethodPost
triggers, err := s.dyn.Resource(httpTrigGVR).Namespace(s.userNS(r)).List(ctx, metav1.ListOptions{})
if err == nil {
for _, trig := range triggers.Items {
refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name")
if refName != name {
continue
}
route, _, _ := unstructured.NestedString(trig.Object, "spec", "relativeurl")
methods, _, _ := unstructured.NestedStringSlice(trig.Object, "spec", "methods")
hasPost := false
hasGet := false
for _, method := range methods {
m := strings.ToUpper(strings.TrimSpace(method))
if m == http.MethodPost {
hasPost = true
}
if m == http.MethodGet {
hasGet = true
}
}
if route != "" {
if !strings.HasPrefix(route, "/") {
route = "/" + route
}
invokeURL = s.routerURL + route
if !hasPost && hasGet {
invokeMethod = http.MethodGet
}
break
}
}
}
start := time.Now()
var invokeBody io.Reader
if invokeMethod == http.MethodPost {
invokeBody = bytes.NewReader(bodyBytes)
}
req, err := http.NewRequestWithContext(ctx, invokeMethod, invokeURL, invokeBody)
if err != nil {
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build invoke request: %v", err))
return
}
if invokeMethod == http.MethodPost {
req.Header.Set("Content-Type", "application/json")
}
if token := s.getRouterToken(); token != "" {
req.Header.Set("Authorization", "Bearer "+token)
}
resp, err := s.http.Do(req)
if err != nil {
if errors.Is(err, context.DeadlineExceeded) {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q timeout after %s: function specialization likely failed (for example, syntax error)", name, invokeTimeout))
return
}
var netErr net.Error
if errors.As(err, &netErr) && netErr.Timeout() {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q timeout after %s: function specialization likely failed (for example, syntax error)", name, invokeTimeout))
return
}
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q: %v", name, err))
return
}
defer resp.Body.Close()
respBody, _ := io.ReadAll(resp.Body)
writeAnyJSON(w, http.StatusOK, map[string]any{
"status": resp.StatusCode,
"latency_ms": time.Since(start).Milliseconds(),
"invoke_url": invokeURL,
"response_raw": string(respBody),
})
}
func (s *server) readSAToken() string {
if s.saTokenPath == "" {
return ""
}
data, err := os.ReadFile(s.saTokenPath)
if err != nil {
return ""
}
return strings.TrimSpace(string(data))
}
func (s *server) getRouterToken() string {
if s.authUser == "" || s.authPass == "" {
return s.readSAToken()
}
s.tokenMu.Lock()
defer s.tokenMu.Unlock()
if s.cachedJWT != "" && time.Now().Before(s.tokenExpAt) {
return s.cachedJWT
}
loginURL := s.routerURL + "/auth/login"
body, _ := json.Marshal(map[string]string{"username": s.authUser, "password": s.authPass})
resp, err := s.http.Post(loginURL, "application/json", bytes.NewReader(body))
if err != nil {
log.Printf("router login failed: %v", err)
return s.readSAToken()
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated {
respBody, _ := io.ReadAll(resp.Body)
log.Printf("router login %d: %s", resp.StatusCode, string(respBody))
return s.readSAToken()
}
var result struct {
AccessToken string `json:"accesstoken"`
}
if err := json.NewDecoder(resp.Body).Decode(&result); err != nil || result.AccessToken == "" {
log.Printf("router login decode error: %v", err)
return s.readSAToken()
}
s.cachedJWT = result.AccessToken
s.tokenExpAt = time.Now().Add(100 * time.Second)
log.Printf("router JWT obtained, expires in 100s")
return s.cachedJWT
}
// --- namespace isolation helpers ---
type ctxKeyNS struct{}
// userNS возвращает namespace пользователя из контекста запроса, или fallback (s.userNS(r)).
func (s *server) userNS(r *http.Request) string {
if ns, ok := r.Context().Value(ctxKeyNS{}).(string); ok && ns != "" {
return ns
}
return s.userNS(r)
}
// namespaceFromJWT декодирует JWT payload (без верификации подписи),
// извлекает claim "sub" и вычисляет namespace: "fission-" + hex(SHA256(sub)[:8]).
func namespaceFromJWT(token string) (string, error) {
parts := strings.SplitN(token, ".", 3)
if len(parts) != 3 {
return "", fmt.Errorf("invalid JWT format")
}
payload := parts[1]
// base64url без padding
switch len(payload) % 4 {
case 2:
payload += "=="
case 3:
payload += "="
}
decoded, err := base64.URLEncoding.DecodeString(payload)
if err != nil {
decoded, err = base64.StdEncoding.DecodeString(payload)
if err != nil {
return "", fmt.Errorf("decode JWT payload: %w", err)
}
}
var claims map[string]any
if err := json.Unmarshal(decoded, &claims); err != nil {
return "", fmt.Errorf("unmarshal JWT claims: %w", err)
}
sub, _ := claims["sub"].(string)
if sub == "" {
return "", fmt.Errorf("JWT missing sub claim")
}
h := sha256.Sum256([]byte(sub))
return "fission-" + hex.EncodeToString(h[:8]), nil
}
// ensureUserNamespace создаёт K8s namespace и shared environments если не существуют.
// addNSToFission динамически добавляет namespace в FISSION_RESOURCE_NAMESPACES у всех Fission deployments.
//
// Зачем это нужно: Fission компоненты смотрят только те namespace-ы, что указаны в
// FISSION_RESOURCE_NAMESPACES. Если не добавить новый namespace — executor/router не будут
// создавать пулы и обрабатывать триггеры, invoke вернёт 404.
//
// Алгоритм:
// 1. Читаем текущее значение FISSION_RESOURCE_NAMESPACES из router deployment.
// 2. Если namespace уже в списке — выходим (idempotent).
// 3. Иначе добавляем namespace к списку и патчим все Fission deployments через StrategicMergePatch.
//
// StrategicMergePatch позволяет обновить только одну env переменную не трогая остальные.
func (s *server) addNSToFission(ctx context.Context, ns string) error {
fissionNS := os.Getenv("FISSION_SYSTEM_NAMESPACE")
if fissionNS == "" {
fissionNS = "fission"
}
fissionDeployments := []string{"router", "executor", "buildermgr", "kubewatcher", "timer"}
// Читаем текущее значение FISSION_RESOURCE_NAMESPACES из router deployment.
// Используем router как источник истины — он первым получает изменения.
routerDep, err := s.dyn.Resource(deploymentGVR).Namespace(fissionNS).Get(ctx, "router", metav1.GetOptions{})
if err != nil {
return fmt.Errorf("get router deployment: %w", err)
}
currentVal := "default" // fallback если переменная не найдена
containerName := "router" // имя контейнера нужно для StrategicMergePatch
containers, _, _ := unstructured.NestedSlice(routerDep.Object, "spec", "template", "spec", "containers")
for _, c := range containers {
cont, ok := c.(map[string]any)
if !ok {
continue
}
// Запоминаем реальное имя контейнера — оно используется как merge key в StrategicMergePatch.
// Без точного имени патч создаст дублирующий контейнер вместо обновления существующего.
if n, ok := cont["name"].(string); ok {
containerName = n
}
envs, _, _ := unstructured.NestedSlice(cont, "env")
for _, e := range envs {
env, ok := e.(map[string]any)
if !ok {
continue
}
if env["name"] == "FISSION_RESOURCE_NAMESPACES" {
if v, ok := env["value"].(string); ok && v != "" {
currentVal = v
}
}
}
break // берём только первый контейнер
}
// Проверяем что namespace ещё не в списке
for _, existing := range strings.Split(currentVal, ",") {
if strings.TrimSpace(existing) == ns {
return nil // уже есть
}
}
newVal := currentVal + "," + ns
// Патчим все Fission deployments одним и тем же значением.
// StrategicMergePatch обновляет только указанные поля (env var), не затрагивая остальные.
// Обычный MergePatch заменил бы весь массив containers — нельзя использовать.
patch := map[string]any{
"spec": map[string]any{
"template": map[string]any{
"spec": map[string]any{
"containers": []any{
map[string]any{
"name": containerName,
"env": []any{
map[string]any{
"name": "FISSION_RESOURCE_NAMESPACES",
"value": newVal,
},
},
},
},
},
},
},
}
patchBytes, err := json.Marshal(patch)
if err != nil {
return fmt.Errorf("marshal patch: %w", err)
}
for _, dep := range fissionDeployments {
// В Fission каждый deployment имеет один контейнер с тем же именем что и deployment.
// Подставляем имя контейнера под конкретный deployment для корректного merge key.
patch["spec"].(map[string]any)["template"].(map[string]any)["spec"].(map[string]any)["containers"].([]any)[0].(map[string]any)["name"] = dep
patchBytes, _ = json.Marshal(patch)
_, patchErr := s.dyn.Resource(deploymentGVR).Namespace(fissionNS).Patch(
ctx, dep, types.StrategicMergePatchType, patchBytes, metav1.PatchOptions{})
if patchErr != nil {
log.Printf("addNSToFission: patch deployment %s: %v", dep, patchErr)
}
}
log.Printf("addNSToFission: added %s, new list: %s", ns, newVal)
return nil
}
func (s *server) ensureUserNamespace(ctx context.Context, ns string) error {
// 1. Создать namespace
nsObj := &unstructured.Unstructured{
Object: map[string]any{
"apiVersion": "v1",
"kind": "Namespace",
"metadata": map[string]any{
"name": ns,
"labels": map[string]any{
"managed-by": "fission-console",
},
},
},
}
_, err := s.dyn.Resource(namespaceGVR).Create(ctx, nsObj, metav1.CreateOptions{})
newlyCreated := err == nil
if err != nil && !apierrors.IsAlreadyExists(err) {
return fmt.Errorf("create namespace %s: %w", ns, err)
}
// 1a. ServiceAccounts для Fission в user namespace.
//
// Fission pool pods (fetcher sidecar) запускаются в user namespace и требуют
// `serviceAccountName: fission-fetcher` в том же namespace. Executor (с SERVICEACCOUNT_CHECK_ENABLED=false)
// не создаёт эти SA автоматически → pool pods падают с "serviceaccount not found".
// Создаём явно при каждом вызове (idempotent через IsAlreadyExists).
saGVR := schema.GroupVersionResource{Group: "", Version: "v1", Resource: "serviceaccounts"}
for _, saName := range []string{"fission-fetcher", "fission-builder"} {
saObj := &unstructured.Unstructured{Object: map[string]any{
"apiVersion": "v1",
"kind": "ServiceAccount",
"metadata": map[string]any{
"name": saName,
"namespace": ns,
},
}}
_, saErr := s.dyn.Resource(saGVR).Namespace(ns).Create(ctx, saObj, metav1.CreateOptions{})
if saErr != nil && !apierrors.IsAlreadyExists(saErr) {
log.Printf("ensureUserNamespace: create SA %s/%s: %v", ns, saName, saErr)
}
}
// 1b. RoleBindings для Fission SA в user namespace.
//
// Проблема: Fission компоненты (executor, router, buildermgr и др.) работают в namespace
// "fission", но при добавлении нового namespace в FISSION_RESOURCE_NAMESPACES они начинают
// туда смотреть (list/watch). По умолчанию у их SA нет прав в чужих namespace-ах →
// "forbidden: cannot list environments.fission.io in namespace X".
//
// Почему cluster-admin, а не admin:
// ClusterRole "admin" не включает custom resource группы (fission.io/*).
// Fission executor при старте пытается создать Role с правами на fission.io/packages,
// и получает "attempting to grant RBAC permissions not currently held" — RBAC escalation
// prevention. ClusterRole "cluster-admin" в контексте RoleBinding (не ClusterRoleBinding)
// даёт полный доступ ТОЛЬКО внутри конкретного namespace — это безопасно.
//
// Операция idempotent: если RoleBinding уже существует — IsAlreadyExists игнорируется.
fissionSAs := []string{"fission-executor", "fission-router", "fission-buildermgr", "fission-kubewatcher", "fission-timer", "fission-fetcher", "fission-builder"}
fissionSysNS := os.Getenv("FISSION_SYSTEM_NAMESPACE")
if fissionSysNS == "" {
fissionSysNS = "fission"
}
rbGVR := schema.GroupVersionResource{Group: "rbac.authorization.k8s.io", Version: "v1", Resource: "rolebindings"}
for _, sa := range fissionSAs {
rbObj := &unstructured.Unstructured{
Object: map[string]any{
"apiVersion": "rbac.authorization.k8s.io/v1",
"kind": "RoleBinding",
"metadata": map[string]any{
"name": "fission-" + sa + "-user-ns",
"namespace": ns,
},
"roleRef": map[string]any{
"apiGroup": "rbac.authorization.k8s.io",
"kind": "ClusterRole",
"name": "cluster-admin", // namespace-scoped через RoleBinding, не ClusterRoleBinding
},
"subjects": []any{
map[string]any{
"kind": "ServiceAccount",
"name": sa,
"namespace": fissionSysNS,
},
},
},
}
_, rbErr := s.dyn.Resource(rbGVR).Namespace(ns).Create(ctx, rbObj, metav1.CreateOptions{})
if rbErr != nil && !apierrors.IsAlreadyExists(rbErr) {
log.Printf("ensureUserNamespace: create rolebinding %s/%s: %v", ns, sa, rbErr)
}
}
// 1c. Регистрируем новый namespace в Fission (FISSION_RESOURCE_NAMESPACES).
// Только при первом создании — повторный патч не нужен, Fission уже знает о namespace.
// addNSToFission читает текущее значение переменной у router-а, добавляет ns и патчит
// все Fission deployments (router, executor, buildermgr, kubewatcher, timer).
if newlyCreated {
if patchErr := s.addNSToFission(ctx, ns); patchErr != nil {
log.Printf("ensureUserNamespace: addNSToFission: %v", patchErr)
}
}
// Environments создаются лениво (lazy) в момент создания первой функции на конкретном языке.
// Это экономит ресурсы: пул подов поднимается только под те языки что реально используются.
return nil
}
func (s *server) validateDeckToken(token, env string) error {
cacheKey := env + ":" + token
if v, ok := s.tokenCache.Load(cacheKey); ok {
if time.Now().Before(v.(time.Time)) {
return nil
}
s.tokenCache.Delete(cacheKey)
}
apiBase, ok := deckAPIs[env]
if !ok {
return fmt.Errorf("unknown env: %s", env)
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
req, err := http.NewRequestWithContext(ctx, http.MethodGet, apiBase+"/index.cfm/instances", nil)
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+token)
resp, err := s.http.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
_, _ = io.ReadAll(resp.Body)
if resp.StatusCode == http.StatusUnauthorized {
return fmt.Errorf("invalid token")
}
s.tokenCache.Store(cacheKey, time.Now().Add(5*time.Minute))
return nil
}
func (s *server) handleAuth(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
writeJSONError(w, http.StatusMethodNotAllowed, "method not allowed")
return
}
var body struct {
Token string `json:"token"`
Env string `json:"env"`
}
if err := json.NewDecoder(r.Body).Decode(&body); err != nil || strings.TrimSpace(body.Token) == "" {
writeJSONError(w, http.StatusBadRequest, "token required")
return
}
env := strings.TrimSpace(strings.ToLower(body.Env))
if _, ok := deckAPIs[env]; !ok {
env = "test"
}
if err := s.validateDeckToken(body.Token, env); err != nil {
writeJSONError(w, http.StatusUnauthorized, "invalid token")
return
}
ns, err := namespaceFromJWT(body.Token)
if err != nil {
log.Printf("handleAuth: namespaceFromJWT: %v", err)
ns = s.ns
}
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second)
defer cancel()
if ensureErr := s.ensureUserNamespace(ctx, ns); ensureErr != nil {
log.Printf("handleAuth: ensureUserNamespace %s: %v", ns, ensureErr)
}
w.Header().Set("Content-Type", "application/json; charset=utf-8")
_ = json.NewEncoder(w).Encode(map[string]any{"ok": true, "env": env, "namespace": ns})
}
func (s *server) handleDeleteFunction(w http.ResponseWriter, r *http.Request, name string) {
ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second)
defer cancel()
var pkgName, envName string
fn, err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).Get(ctx, name, metav1.GetOptions{})
if err != nil {
if apierrors.IsNotFound(err) {
writeJSONError(w, http.StatusNotFound, fmt.Sprintf("function %q not found", name))
return
}
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get function %q: %v", name, err))
return
}
pkgName, _, _ = unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name")
envName, _, _ = unstructured.NestedString(fn.Object, "spec", "environment", "name")
triggers, err := s.dyn.Resource(httpTrigGVR).Namespace(s.userNS(r)).List(ctx, metav1.ListOptions{})
if err == nil {
for _, trig := range triggers.Items {
refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name")
if refName == name {
_ = s.dyn.Resource(httpTrigGVR).Namespace(s.userNS(r)).Delete(ctx, trig.GetName(), metav1.DeleteOptions{})
}
}
}
if err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).Delete(ctx, name, metav1.DeleteOptions{}); err != nil && !apierrors.IsNotFound(err) {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("delete function %q: %v", name, err))
return
}
if pkgName != "" {
if err := s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Delete(ctx, pkgName, metav1.DeleteOptions{}); err != nil && !apierrors.IsNotFound(err) {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("delete package %q: %v", pkgName, err))
return
}
}
// Если environment больше не используется ни одной функцией — удаляем его.
// Fission увидит удаление Environment CRD и убьёт pool deployment → поды умирают.
if envName != "" {
cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cleanupCancel()
s.cleanupEnvironmentIfUnused(cleanupCtx, s.userNS(r), envName)
}
writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name, "package": pkgName})
}
func buildConfig(kubeconfig string) (*rest.Config, error) {
if kubeconfig != "" {
cfg, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
if err == nil {
return cfg, nil
}
return nil, fmt.Errorf("kubeconfig %s: %w", kubeconfig, err)
}
cfg, err := rest.InClusterConfig()
if err == nil {
return cfg, nil
}
loadingRules := &clientcmd.ClientConfigLoadingRules{}
clientCfg := clientcmd.NewNonInteractiveDeferredLoadingClientConfig(loadingRules, &clientcmd.ConfigOverrides{})
return clientCfg.ClientConfig()
}
func (s *server) handleList(gvr schema.GroupVersionResource) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
return
}
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second)
defer cancel()
list, err := s.dyn.Resource(gvr).Namespace(s.userNS(r)).List(ctx, metav1.ListOptions{})
if err != nil {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("list %s: %v", gvr.Resource, err))
return
}
writeJSON(w, http.StatusOK, list.Items)
}
}
func writeJSON(w http.ResponseWriter, status int, data []unstructured.Unstructured) {
w.Header().Set("Content-Type", "application/json; charset=utf-8")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(data)
}
func writeAnyJSON(w http.ResponseWriter, status int, data any) {
w.Header().Set("Content-Type", "application/json; charset=utf-8")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(data)
}
func writeJSONError(w http.ResponseWriter, status int, msg string) {
w.Header().Set("Content-Type", "application/json; charset=utf-8")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(map[string]any{"error": msg})
}
func logRequests(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
log.Printf("%s %s", r.Method, r.URL.Path)
next.ServeHTTP(w, r)
})
}
func withCORS(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Access-Control-Allow-Origin", "*")
w.Header().Set("Access-Control-Allow-Methods", "GET,POST,PUT,PATCH,DELETE,OPTIONS")
w.Header().Set("Access-Control-Allow-Headers", "Content-Type, Authorization, X-Auth-Token, X-Auth-Env")
if r.Method == http.MethodOptions {
w.WriteHeader(http.StatusNoContent)
return
}
next.ServeHTTP(w, r)
})
}
func withSecurityHeaders(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("X-Content-Type-Options", "nosniff")
w.Header().Set("X-Frame-Options", "DENY")
w.Header().Set("Referrer-Policy", "strict-origin-when-cross-origin")
w.Header().Set("Permissions-Policy", "camera=(), microphone=(), geolocation=()")
w.Header().Set("Content-Security-Policy", "default-src 'self'; script-src 'self' 'unsafe-inline'; style-src 'self' 'unsafe-inline'; img-src 'self' data:; connect-src 'self'; font-src 'self' data:; object-src 'none'; frame-ancestors 'none'; base-uri 'self'; form-action 'self'; upgrade-insecure-requests; block-all-mixed-content")
next.ServeHTTP(w, r)
})
}
func envDefault(key, fallback string) string {
if v := strings.TrimSpace(os.Getenv(key)); v != "" {
return v
}
return fallback
}
func envDurationDefault(key string, fallback time.Duration) time.Duration {
raw := strings.TrimSpace(os.Getenv(key))
if raw == "" {
return fallback
}
d, err := time.ParseDuration(raw)
if err != nil || d <= 0 {
log.Printf("invalid duration for %s=%q, using default %s", key, raw, fallback)
return fallback
}
return d
}
func normalizeMethods(in []string) []string {
if len(in) == 0 {
return []string{"GET"}
}
out := make([]string, 0, len(in))
seen := map[string]bool{}
for _, method := range in {
m := strings.ToUpper(strings.TrimSpace(method))
if m == "" || seen[m] {
continue
}
seen[m] = true
out = append(out, m)
}
if len(out) == 0 {
return []string{"GET"}
}
return out
}
func (s *server) extractPackageSourceCode(ctx context.Context, pkg *unstructured.Unstructured) string {
literalPaths := [][]string{
{"spec", "source", "literal"},
{"spec", "deployment", "literal"},
}
for _, p := range literalPaths {
literal, found, _ := unstructured.NestedString(pkg.Object, p...)
if !found || strings.TrimSpace(literal) == "" {
continue
}
if decodedCode, decErr := decodeLiteralToSource(literal); decErr == nil && strings.TrimSpace(decodedCode) != "" {
return decodedCode
}
}
urlPaths := [][]string{
{"spec", "source", "url"},
{"spec", "deployment", "url"},
}
for _, p := range urlPaths {
urlValue, found, _ := unstructured.NestedString(pkg.Object, p...)
if !found || strings.TrimSpace(urlValue) == "" {
continue
}
archiveBytes, fetchErr := s.fetchPackageArchive(ctx, urlValue)
if fetchErr != nil {
continue
}
decodedCode, decErr := decodeArchiveBytesToSource(archiveBytes)
if decErr == nil && strings.TrimSpace(decodedCode) != "" {
return decodedCode
}
}
return ""
}
func (s *server) fetchPackageArchive(ctx context.Context, archiveURL string) ([]byte, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, archiveURL, nil)
if err != nil {
return nil, err
}
resp, err := s.http.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return nil, fmt.Errorf("archive request failed: %s", resp.Status)
}
return io.ReadAll(resp.Body)
}
func decodeLiteralToSource(literal string) (string, error) {
decoded, err := base64.StdEncoding.DecodeString(literal)
if err != nil {
return "", err
}
return decodeArchiveBytesToSource(decoded)
}
func decodeArchiveBytesToSource(decoded []byte) (string, error) {
if len(decoded) == 0 {
return "", fmt.Errorf("empty payload")
}
if utf8.Valid(decoded) {
return string(decoded), nil
}
if len(decoded) >= 4 && bytes.Equal(decoded[:4], []byte{'P', 'K', 3, 4}) {
if src, zipErr := decodeZipSource(decoded); zipErr == nil {
return src, nil
}
}
return "", fmt.Errorf("payload does not contain utf-8 source")
}
func decodeZipSource(zipBytes []byte) (string, error) {
reader, err := zip.NewReader(bytes.NewReader(zipBytes), int64(len(zipBytes)))
if err != nil {
return "", err
}
preferred := []string{"main.py", "main.js", "main.go", "handler.go", "handler.js", "handler.py"}
for _, name := range preferred {
for _, file := range reader.File {
if strings.EqualFold(file.Name, name) {
content, readErr := readZipFile(file)
if readErr != nil {
return "", readErr
}
if utf8.Valid(content) {
return string(content), nil
}
}
}
}
files := make([]*zip.File, 0, len(reader.File))
for _, file := range reader.File {
if file.FileInfo().IsDir() {
continue
}
files = append(files, file)
}
sort.Slice(files, func(i, j int) bool {
return files[i].Name < files[j].Name
})
for _, file := range files {
content, readErr := readZipFile(file)
if readErr != nil {
continue
}
if utf8.Valid(content) {
return string(content), nil
}
}
return "", fmt.Errorf("zip archive does not contain utf-8 source files")
}
func readZipFile(file *zip.File) ([]byte, error) {
rc, err := file.Open()
if err != nil {
return nil, err
}
defer rc.Close()
return io.ReadAll(rc)
}