Files
fission-console/console/main.go
T

2431 lines
83 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"
"go/parser"
"go/token"
"io"
"log"
"net"
"net/http"
"os"
"os/exec"
"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",
}
// nsInflightEnsure — состояние in-flight вызова ensureUserNamespace.
// Все параллельные горутины ждут close(done), затем читают err.
type nsInflightEnsure struct {
done chan struct{}
err error
}
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
// --- ai/ask feature (удалить блок целиком чтобы выкосить) ---
llmURL string // FISSION_LLM_URL
llmKey string // FISSION_LLM_KEY
// --- end ai/ask feature ---
tokenMu sync.Mutex
cachedJWT string
tokenExpAt time.Time
tokenCache sync.Map
// nsReconcileCh — сигнал для немедленного запуска NS reconciler.
// Буферизирован на 1: несколько сигналов схлопываются в один запуск.
nsReconcileCh chan struct{}
// ensuredNS — кэш namespace-ов для которых уже отработал ensureUserNamespace.
// ensuredNSMu защищает ensuredNS и ensuredNSInFlight.
// ensuredNSInFlight — ожидание: если namespace создаётся прямо сейчас, параллельные
// запросы ждут завершения (ручной singleflight без внешних зависимостей).
ensuredNSMu sync.Mutex
ensuredNS map[string]struct{}
ensuredNSInFlight map[string]*nsInflightEnsure
// nsSemaphore ограничивает параллелизм ensureUserNamespace — не более 3 одновременно.
// Без него 10 новых пользователей генерируют 140 K8s API calls одновременно → throttle → 504.
nsSemaphore chan struct{}
}
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: "naeel/go-builder-fast:v1"},
"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",
// --- ai/ask feature ---
llmURL: envDefault("FISSION_LLM_URL", "https://api.aillm.ru"),
llmKey: os.Getenv("FISSION_LLM_KEY"),
// --- end ai/ask feature ---
nsReconcileCh: make(chan struct{}, 1),
ensuredNS: make(map[string]struct{}),
ensuredNSInFlight: make(map[string]*nsInflightEnsure),
nsSemaphore: make(chan struct{}, 3),
}
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 или X-Auth-Token задают sub → разные namespace-ы для тестирования.
sub := strings.TrimSpace(r.Header.Get("X-Test-Sub"))
if sub == "" {
sub = strings.TrimSpace(r.Header.Get("X-Auth-Token"))
}
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)
// Гарантируем что namespace + RBAC + quota + netpol существуют.
// ensureUserNS реализует singleflight + кэш + семафор параллелизма.
if ensureErr := s.ensureUserNS(ctx, ns); ensureErr != nil {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure namespace: %v", ensureErr))
return
}
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)))
mux.HandleFunc("/console/api/ns/status", auth(s.handleNSStatus))
mux.HandleFunc("/console/api/ai/check", auth(s.handleAICheck))
// --- ai/ask feature (удалить строку чтобы выкосить роут) ---
mux.HandleFunc("/console/api/ai/ask", auth(s.handleAIAsk))
// --- end ai/ask feature ---
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))
s.startNSReconciler(envDurationDefault("NS_RECONCILE_INTERVAL", 2*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 через кэш+singleflight+семафор.
nsCtx, nsCancel := context.WithTimeout(r.Context(), 60*time.Second)
defer nsCancel()
if err := s.ensureUserNS(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)
}
}
// Сигналим reconciler — он проверит все namespace'ы и почистит список.
select {
case s.nsReconcileCh <- struct{}{}:
default:
}
// Сканируем 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 если не существуют.
// После создания сигналит nsReconciler который синхронизирует FISSION_RESOURCE_NAMESPACES.
//
// Зачем reconciler, а не прямой патч:
// - Прямой патч на горячем пути → rolling restart всех Fission deployments при каждом новом юзере
// - Race condition: 100 юзеров одновременно → read-modify-write без мьютекса → namespace'ы теряются
// - Reconciler батчит изменения, нет race condition, нет лишних restarts
func (s *server) addNSToFission(_ context.Context, _ string) error { return nil } // заменено reconciler'ом
// startNSReconciler запускает фоновый reconciler FISSION_RESOURCE_NAMESPACES.
//
// Проблема при scale:
// - addNSToFission на горячем пути (каждый новый юзер) → патч 5 деплоев → rolling restart → cold start для всех
// - Без мьютекса: 100 юзеров одновременно → read-modify-write race → namespace'ы теряются
// - Ручное kubectl delete ns → namespace остаётся в переменной вечно → executor флудит RBAC ошибками
//
// Решение: единственный goroutine с debounce-каналом.
// - Запускается по таймеру (каждые NS_RECONCILE_INTERVAL) или немедленно через nsReconcileCh
// - 1000 юзеров создают namespace'ы одновременно → 1 патч вместо 5000
// - Нет race condition (один goroutine, один writer)
// - Автоматически чистит ghost namespace'ы (удалённые kubectl delete ns или reaper'ом)
func (s *server) startNSReconciler(interval time.Duration) {
go func() {
ticker := time.NewTicker(interval)
defer ticker.Stop()
log.Printf("nsReconciler: started, interval=%v", interval)
for {
select {
case <-ticker.C:
s.reconcileNSList()
case <-s.nsReconcileCh:
// Немедленный запуск (новый namespace или удаление функции).
// Дренируем канал чтобы не запускаться дважды подряд.
s.reconcileNSList()
drain:
for {
select {
case <-s.nsReconcileCh:
default:
break drain
}
}
}
}
}()
}
// reconcileNSList синхронизирует FISSION_RESOURCE_NAMESPACES с реально существующими namespace'ами.
//
// Алгоритм:
// 1. Читает все namespace'ы с меткой managed-by=fission-console из k8s (источник истины)
// 2. Читает текущий FISSION_RESOURCE_NAMESPACES из router deployment
// 3. Вычисляет desired = "default" + активные namespace'ы (существующие и непустые)
// 4. Если desired == current → ничего не делает (нет rolling restart!)
// 5. Если разница → один патч всех Fission deployments
func (s *server) reconcileNSList() {
ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
defer cancel()
fissionNS := os.Getenv("FISSION_SYSTEM_NAMESPACE")
if fissionNS == "" {
fissionNS = "fission"
}
// Шаг 1: реальные namespace'ы с нашей меткой
nsList, err := s.dyn.Resource(namespaceGVR).List(ctx, metav1.ListOptions{
LabelSelector: "managed-by=fission-console",
})
if err != nil {
log.Printf("nsReconciler: list namespaces: %v", err)
return
}
// Шаг 2: фильтруем — берём только те что Active.
// Terminating namespace'ы убираем из списка (они уже умирают).
desired := map[string]struct{}{"default": {}}
for _, ns := range nsList.Items {
phase, _, _ := unstructured.NestedString(ns.Object, "status", "phase")
if phase == "Active" {
desired[ns.GetName()] = struct{}{}
}
}
// Шаг 3: текущее значение из router (источник истины)
routerDep, err := s.dyn.Resource(deploymentGVR).Namespace(fissionNS).Get(ctx, "router", metav1.GetOptions{})
if err != nil {
log.Printf("nsReconciler: get router: %v", err)
return
}
currentVal := "default"
containers, _, _ := unstructured.NestedSlice(routerDep.Object, "spec", "template", "spec", "containers")
for _, c := range containers {
cont, ok := c.(map[string]any)
if !ok {
continue
}
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
}
// Шаг 4: сравниваем current с desired
currentSet := map[string]struct{}{}
for _, p := range strings.Split(currentVal, ",") {
if t := strings.TrimSpace(p); t != "" {
currentSet[t] = struct{}{}
}
}
// Проверяем симметричную разницу
same := len(currentSet) == len(desired)
if same {
for k := range desired {
if _, ok := currentSet[k]; !ok {
same = false
break
}
}
}
if same {
return // ничего менять не нужно — нет патча, нет rolling restart
}
// Шаг 5: патчим все Fission deployments
parts := make([]string, 0, len(desired))
for ns := range desired {
parts = append(parts, ns)
}
sort.Strings(parts)
newVal := strings.Join(parts, ",")
fissionDeployments := []string{"router", "executor", "buildermgr", "kubewatcher", "timer"}
patch := map[string]any{
"spec": map[string]any{
"template": map[string]any{
"spec": map[string]any{
"containers": []any{
map[string]any{
"name": "",
"env": []any{
map[string]any{
"name": "FISSION_RESOURCE_NAMESPACES",
"value": newVal,
},
},
},
},
},
},
},
}
for _, dep := range fissionDeployments {
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)
if _, pErr := s.dyn.Resource(deploymentGVR).Namespace(fissionNS).Patch(
ctx, dep, types.StrategicMergePatchType, patchBytes, metav1.PatchOptions{}); pErr != nil {
log.Printf("nsReconciler: patch %s: %v", dep, pErr)
}
}
log.Printf("nsReconciler: synced FISSION_RESOURCE_NAMESPACES: %s → %s", currentVal, newVal)
}
// ensureUserNS — единая точка входа для гарантии существования пользовательского namespace.
// Реализует singleflight + in-memory кэш + семафор параллелизма.
//
// Singleflight: если один goroutine уже создаёт namespace ns — остальные ждут его результата
// вместо того чтобы запускать параллельные K8s API calls (вызывало throttle и 504).
//
// Кэш: если namespace уже создан в этом запуске процесса — быстрый путь без K8s calls.
//
// Семафор (3 слота): не более 3 namespace-ов создаются одновременно.
// 10 новых пользователей × 14 K8s calls = 140 calls без семафора → throttle → 60s+ → 504.
// С семафором: 3 batch-а по 14 calls → ~3 × 5s = 15s total, все укладываются в timeout.
func (s *server) ensureUserNS(ctx context.Context, ns string) error {
s.ensuredNSMu.Lock()
if _, ok := s.ensuredNS[ns]; ok {
// Быстрый путь: уже создан в этой жизни процесса.
s.ensuredNSMu.Unlock()
return nil
}
if inflight, ok := s.ensuredNSInFlight[ns]; ok {
// Кто-то уже создаёт — ждём его результата.
s.ensuredNSMu.Unlock()
select {
case <-inflight.done:
return inflight.err
case <-ctx.Done():
return ctx.Err()
}
}
// Мы первые для этого namespace.
inflight := &nsInflightEnsure{done: make(chan struct{})}
s.ensuredNSInFlight[ns] = inflight
s.ensuredNSMu.Unlock()
// Берём слот семафора — ограничиваем параллелизм.
select {
case s.nsSemaphore <- struct{}{}:
case <-ctx.Done():
s.ensuredNSMu.Lock()
delete(s.ensuredNSInFlight, ns)
s.ensuredNSMu.Unlock()
inflight.err = ctx.Err()
close(inflight.done)
return ctx.Err()
}
ensureCtx, ensureCancel := context.WithTimeout(ctx, 60*time.Second)
inflight.err = s.ensureUserNamespace(ensureCtx, ns)
ensureCancel()
<-s.nsSemaphore // освобождаем слот
s.ensuredNSMu.Lock()
delete(s.ensuredNSInFlight, ns)
if inflight.err == nil {
s.ensuredNS[ns] = struct{}{}
}
s.ensuredNSMu.Unlock()
close(inflight.done)
return inflight.err
}
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 игнорируется.
// Важно: pool pod запускается с serviceAccountName=fission-fetcher в самом user namespace,
// а не в system namespace "fission". Поэтому для fetcher/builder нужны ещё локальные bindings.
type rbSubject struct {
name string
namespace string
binding string
}
fissionSAs := []rbSubject{
{name: "fission-executor", binding: "fission-executor-user-ns"},
{name: "fission-router", binding: "fission-router-user-ns"},
{name: "fission-buildermgr", binding: "fission-buildermgr-user-ns"},
{name: "fission-kubewatcher", binding: "fission-kubewatcher-user-ns"},
{name: "fission-timer", binding: "fission-timer-user-ns"},
{name: "fission-fetcher", binding: "fission-fetcher-system-user-ns"},
{name: "fission-builder", binding: "fission-builder-system-user-ns"},
{name: "fission-fetcher", namespace: ns, binding: "fission-fetcher-local-user-ns"},
{name: "fission-builder", namespace: ns, binding: "fission-builder-local-user-ns"},
}
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 {
subjectNS := sa.namespace
if subjectNS == "" {
subjectNS = fissionSysNS
}
rbObj := &unstructured.Unstructured{
Object: map[string]any{
"apiVersion": "rbac.authorization.k8s.io/v1",
"kind": "RoleBinding",
"metadata": map[string]any{
"name": sa.binding,
"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.name,
"namespace": subjectNS,
},
},
},
}
_, 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@%s: %v", ns, sa.name, subjectNS, rbErr)
}
}
// 1c. ResourceQuota — ограничиваем сколько ресурсов может потребить один пользователь.
// Без этого одна функция может исчерпать CPU/RAM всего кластера.
// Параметры вынесены в env vars для гибкой настройки.
quotaGVR := schema.GroupVersionResource{Group: "", Version: "v1", Resource: "resourcequotas"}
quotaObj := &unstructured.Unstructured{Object: map[string]any{
"apiVersion": "v1",
"kind": "ResourceQuota",
"metadata": map[string]any{
"name": "user-quota",
"namespace": ns,
},
"spec": map[string]any{
"hard": map[string]any{
"requests.cpu": envDefault("QUOTA_REQ_CPU", "1"),
"requests.memory": envDefault("QUOTA_REQ_MEM", "1Gi"),
"limits.cpu": envDefault("QUOTA_LIM_CPU", "4"),
"limits.memory": envDefault("QUOTA_LIM_MEM", "4Gi"),
"pods": envDefault("QUOTA_PODS", "30"),
"count/functions.fission.io": envDefault("QUOTA_FUNCTIONS", "20"),
"count/packages.fission.io": envDefault("QUOTA_PACKAGES", "40"),
"count/httptriggers.fission.io": envDefault("QUOTA_HTTPTRIGGERS", "20"),
},
},
}}
_, quotaErr := s.dyn.Resource(quotaGVR).Namespace(ns).Create(ctx, quotaObj, metav1.CreateOptions{})
if quotaErr != nil && !apierrors.IsAlreadyExists(quotaErr) {
log.Printf("ensureUserNamespace: create ResourceQuota %s: %v", ns, quotaErr)
}
// 1d. LimitRange — дефолтные лимиты на контейнер, чтобы поды без явных limits не были unbounded.
limitRangeGVR := schema.GroupVersionResource{Group: "", Version: "v1", Resource: "limitranges"}
limitRangeObj := &unstructured.Unstructured{Object: map[string]any{
"apiVersion": "v1",
"kind": "LimitRange",
"metadata": map[string]any{
"name": "user-limits",
"namespace": ns,
},
"spec": map[string]any{
"limits": []any{
map[string]any{
"type": "Container",
"default": map[string]any{
"cpu": envDefault("LIMIT_DEFAULT_CPU", "500m"),
"memory": envDefault("LIMIT_DEFAULT_MEM", "256Mi"),
},
"defaultRequest": map[string]any{
"cpu": envDefault("LIMIT_REQ_CPU", "50m"),
"memory": envDefault("LIMIT_REQ_MEM", "64Mi"),
},
"max": map[string]any{
"cpu": envDefault("LIMIT_MAX_CPU", "2"),
"memory": envDefault("LIMIT_MAX_MEM", "1Gi"),
},
},
},
},
}}
_, lrErr := s.dyn.Resource(limitRangeGVR).Namespace(ns).Create(ctx, limitRangeObj, metav1.CreateOptions{})
if lrErr != nil && !apierrors.IsAlreadyExists(lrErr) {
log.Printf("ensureUserNamespace: create LimitRange %s: %v", ns, lrErr)
}
// 1e. NetworkPolicy — запрещаем входящий трафик из других user namespace.
// Разрешаем:
// - трафик внутри самого namespace (pod → pod в том же ns)
// - трафик из fission core namespace (router → function pod)
// - трафик из kube-system (kubelet health checks, DNS)
// Запрещаем:
// - трафик от подов других user namespace (межтенантная изоляция)
fissionSysNSForNetpol := os.Getenv("FISSION_SYSTEM_NAMESPACE")
if fissionSysNSForNetpol == "" {
fissionSysNSForNetpol = "fission"
}
netpolGVR := schema.GroupVersionResource{Group: "networking.k8s.io", Version: "v1", Resource: "networkpolicies"}
netpolObj := &unstructured.Unstructured{Object: map[string]any{
"apiVersion": "networking.k8s.io/v1",
"kind": "NetworkPolicy",
"metadata": map[string]any{
"name": "deny-cross-tenant",
"namespace": ns,
},
"spec": map[string]any{
"podSelector": map[string]any{}, // применяется ко всем подам namespace
"policyTypes": []any{"Ingress"},
"ingress": []any{
// Разрешаем трафик внутри namespace
map[string]any{
"from": []any{
map[string]any{
"podSelector": map[string]any{},
},
},
},
// Разрешаем трафик из fission core namespace (router, executor)
map[string]any{
"from": []any{
map[string]any{
"namespaceSelector": map[string]any{
"matchLabels": map[string]any{
"kubernetes.io/metadata.name": fissionSysNSForNetpol,
},
},
},
},
},
// Разрешаем трафик из kube-system (DNS, health checks)
map[string]any{
"from": []any{
map[string]any{
"namespaceSelector": map[string]any{
"matchLabels": map[string]any{
"kubernetes.io/metadata.name": "kube-system",
},
},
},
},
},
},
},
}}
_, npErr := s.dyn.Resource(netpolGVR).Namespace(ns).Create(ctx, netpolObj, metav1.CreateOptions{})
if npErr != nil && !apierrors.IsAlreadyExists(npErr) {
log.Printf("ensureUserNamespace: create NetworkPolicy %s: %v", ns, npErr)
}
// 1f. Сигналим NS reconciler что появился новый namespace.
// Reconciler сам синхронизирует FISSION_RESOURCE_NAMESPACES — без race condition и rolling restarts.
if newlyCreated {
select {
case s.nsReconcileCh <- struct{}{}:
default: // уже есть сигнал в буфере — не блокируем
}
}
// Environments создаются лениво (lazy) в момент создания первой функции на конкретном языке.
// Это экономит ресурсы: пул подов поднимается только под те языки что реально используются.
return nil
}
// handleNSStatus возвращает статус инициализации пользовательского namespace.
// Используется UI для отображения прогресса при первом входе.
func (s *server) handleNSStatus(w http.ResponseWriter, r *http.Request) {
ns := s.userNS(r)
ctx, cancel := context.WithTimeout(r.Context(), 15*time.Second)
defer cancel()
fissionNS := os.Getenv("FISSION_SYSTEM_NAMESPACE")
if fissionNS == "" {
fissionNS = "fission"
}
type stageInfo struct {
Name string `json:"name"`
Done bool `json:"done"`
}
stages := []stageInfo{
{Name: "Создание пространства имён"},
{Name: "Регистрация в Fission"},
{Name: "Прогрев окружений"},
}
// Stage 1: NS существует и Active (ensureUserNS уже вызван auth middleware)
nsObj, err := s.dyn.Resource(namespaceGVR).Get(ctx, ns, metav1.GetOptions{})
if err == nil {
phase, _, _ := unstructured.NestedString(nsObj.Object, "status", "phase")
stages[0].Done = phase == "Active"
}
// Stage 2: NS в FISSION_RESOURCE_NAMESPACES executor deployment
if stages[0].Done {
execDep, err2 := s.dyn.Resource(deploymentGVR).Namespace(fissionNS).Get(ctx, "executor", metav1.GetOptions{})
if err2 == nil {
containers, _, _ := unstructured.NestedSlice(execDep.Object, "spec", "template", "spec", "containers")
for _, c := range containers {
cont, ok := c.(map[string]any)
if !ok {
continue
}
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 {
for _, p := range strings.Split(v, ",") {
if strings.TrimSpace(p) == ns {
stages[1].Done = true
}
}
}
}
}
break
}
}
}
// Stage 3: executor pod Running + Ready
if stages[1].Done {
podGVR := schema.GroupVersionResource{Group: "", Version: "v1", Resource: "pods"}
podList, err3 := s.dyn.Resource(podGVR).Namespace(fissionNS).List(ctx, metav1.ListOptions{
LabelSelector: "svc=executor",
})
if err3 == nil {
for _, pod := range podList.Items {
phase, _, _ := unstructured.NestedString(pod.Object, "status", "phase")
if phase != "Running" {
continue
}
conditions, _, _ := unstructured.NestedSlice(pod.Object, "status", "conditions")
for _, c := range conditions {
cond, ok := c.(map[string]any)
if !ok {
continue
}
if cond["type"] == "Ready" && cond["status"] == "True" {
stages[2].Done = true
}
}
}
}
}
ready := stages[0].Done && stages[1].Done && stages[2].Done
writeAnyJSON(w, http.StatusOK, map[string]any{
"ready": ready,
"stages": stages,
})
}
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"
}
var ns string
if s.testMode {
// testMode: токен — это sub (email), Deck не вызывается
if !strings.Contains(body.Token, "@") {
writeJSONError(w, http.StatusUnauthorized, "invalid token")
return
}
ns = func() string {
h32 := sha256.Sum256([]byte(body.Token))
return "fission-" + hex.EncodeToString(h32[:8])
}()
} else {
if err := s.validateDeckToken(body.Token, env); err != nil {
writeJSONError(w, http.StatusUnauthorized, "invalid token")
return
}
var err error
ns, err = namespaceFromJWT(body.Token)
if err != nil {
log.Printf("handleAuth: namespaceFromJWT: %v", err)
ns = s.ns
}
}
ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second)
defer cancel()
if ensureErr := s.ensureUserNS(ctx, ns); ensureErr != nil {
log.Printf("handleAuth: ensureUserNS %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 больше не используется ни одной функцией — удаляем его.
// Если 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)
}
// Сигналим reconciler: он проверит все namespace'ы и уберёт пустые из FISSION_RESOURCE_NAMESPACES.
// Не делаем это на горячем пути — reconciler батчит изменения без race condition и rolling restarts.
select {
case s.nsReconcileCh <- struct{}{}:
default:
}
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)
}
// goParse проверяет синтаксис Go-кода через go/parser.
func goParse(code string) (*token.FileSet, error) {
fset := token.NewFileSet()
_, err := parser.ParseFile(fset, "code.go", code, parser.AllErrors)
if err != nil {
return nil, err
}
return fset, nil
}
// handleAICheck проверяет синтаксис кода через реальный линтер языка.
// POST /console/api/ai/check
// Body: {"language":"python","code":"..."}
// Response: {"ok":true,"result":"..."}
func (s *server) handleAICheck(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
writeJSONError(w, http.StatusMethodNotAllowed, "method not allowed")
return
}
var req struct {
Language string `json:"language"`
Code string `json:"code"`
}
if err := json.NewDecoder(io.LimitReader(r.Body, 64*1024)).Decode(&req); err != nil {
writeJSONError(w, http.StatusBadRequest, "invalid JSON: "+err.Error())
return
}
if strings.TrimSpace(req.Code) == "" {
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]any{"ok": false, "result": "Код пустой."})
return
}
// Определяем команду и расширение файла по языку
type linterCfg struct {
ext string
cmd []string
}
langs := map[string]linterCfg{
"nodejs": {ext: ".js", cmd: []string{"node", "--check"}},
"python": {ext: ".py", cmd: []string{"python3", "-m", "py_compile"}},
"ruby": {ext: ".rb", cmd: []string{"ruby", "-c"}},
"php": {ext: ".php", cmd: []string{"php", "-l"}},
"perl": {ext: ".pl", cmd: []string{"perl", "-c"}},
"go": {ext: ".go", cmd: nil}, // go проверяем через go/parser
}
cfg, ok := langs[req.Language]
if !ok {
writeJSONError(w, http.StatusBadRequest, "unsupported language: "+req.Language)
return
}
var isOK bool
var result string
if req.Language == "go" {
// Go: используем go/parser прямо в процессе — без внешних команд
_, parseErr := goParse(req.Code)
if parseErr == nil {
isOK = true
result = "✅ Синтаксис корректен."
} else {
isOK = false
result = parseErr.Error()
}
} else {
// Записываем код во временный файл
tmpf, err := os.CreateTemp("", "fission-lint-*"+cfg.ext)
if err != nil {
writeJSONError(w, http.StatusInternalServerError, "tmp file: "+err.Error())
return
}
defer os.Remove(tmpf.Name())
if _, err := tmpf.WriteString(req.Code); err != nil {
tmpf.Close()
writeJSONError(w, http.StatusInternalServerError, "write tmp: "+err.Error())
return
}
tmpf.Close()
args := append(cfg.cmd, tmpf.Name())
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second)
defer cancel()
//nolint:gosec — cfg.cmd содержит только захардкоженные команды из langsMap
out, err := exec.CommandContext(ctx, args[0], args[1:]...).CombinedOutput()
outStr := strings.TrimSpace(string(out))
// Убираем путь к tmp-файлу из вывода — юзеру незачем его видеть
outStr = strings.ReplaceAll(outStr, tmpf.Name(), "<code>")
if err == nil {
isOK = true
result = "✅ Синтаксис корректен."
} else {
isOK = false
if outStr != "" {
result = outStr
} else {
result = "Ошибка синтаксиса (линтер вернул код " + fmt.Sprintf("%v", err) + ")"
}
}
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]any{
"ok": isOK,
"result": result,
})
}
// =============================================================================
// ai/ask feature — удалить весь блок до "end ai/ask feature" чтобы выкосить
// =============================================================================
// handleAIAsk отвечает на произвольный вопрос пользователя через LLM.
// POST /console/api/ai/ask
// Body: {"question":"..."}
// Response: {"answer":"..."}
func (s *server) handleAIAsk(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
writeJSONError(w, http.StatusMethodNotAllowed, "method not allowed")
return
}
if s.llmKey == "" {
writeJSONError(w, http.StatusServiceUnavailable, "AI-ассистент не настроен (FISSION_LLM_KEY не задан)")
return
}
var req struct {
Question string `json:"question"`
}
if err := json.NewDecoder(io.LimitReader(r.Body, 4*1024)).Decode(&req); err != nil {
writeJSONError(w, http.StatusBadRequest, "invalid JSON: "+err.Error())
return
}
if strings.TrimSpace(req.Question) == "" {
writeJSONError(w, http.StatusBadRequest, "question is required")
return
}
body, _ := json.Marshal(map[string]any{
"model": "gpt-oss-120b",
"messages": []map[string]string{
{
"role": "system",
"content": "Ты умный ассистент. Отвечай кратко и по делу.",
},
{"role": "user", "content": req.Question},
},
"max_tokens": 1024,
})
ctx, cancel := context.WithTimeout(r.Context(), 30*time.Second)
defer cancel()
llmReq, err := http.NewRequestWithContext(ctx, http.MethodPost,
strings.TrimRight(s.llmURL, "/")+"/chat/completions",
bytes.NewReader(body),
)
if err != nil {
writeJSONError(w, http.StatusInternalServerError, err.Error())
return
}
llmReq.Header.Set("Content-Type", "application/json")
llmReq.Header.Set("Authorization", "Bearer "+s.llmKey)
resp, err := s.http.Do(llmReq)
if err != nil {
writeJSONError(w, http.StatusBadGateway, "LLM недоступен: "+err.Error())
return
}
defer resp.Body.Close()
var llmResp struct {
Choices []struct {
Message struct {
Content string `json:"content"`
} `json:"message"`
} `json:"choices"`
Error *struct {
Message string `json:"message"`
} `json:"error"`
}
if err := json.NewDecoder(io.LimitReader(resp.Body, 64*1024)).Decode(&llmResp); err != nil {
writeJSONError(w, http.StatusBadGateway, "parse error: "+err.Error())
return
}
if llmResp.Error != nil {
writeJSONError(w, http.StatusBadGateway, llmResp.Error.Message)
return
}
if len(llmResp.Choices) == 0 {
writeJSONError(w, http.StatusBadGateway, "empty response from LLM")
return
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]any{
"answer": strings.TrimSpace(llmResp.Choices[0].Message.Content),
})
}
// =============================================================================
// end ai/ask feature
// =============================================================================