v1.3.27: fix Help modal — add -L flag, demo-login docs

This commit is contained in:
Naeel
2026-05-03 06:46:27 +03:00
parent d7352532dc
commit df74a69fde
37 changed files with 1179 additions and 2497 deletions
+45 -110
View File
@@ -2,19 +2,20 @@ package api
import (
"context"
"crypto/sha256"
"encoding/base64"
"encoding/hex"
"encoding/json"
"fmt"
"net/http"
"strings"
"fission-console/internal/auth"
)
// ctxKeyNS — ключ для хранения namespace пользователя в context.Context.
// Использует приватный тип чтобы избежать коллизий с ключами из других пакетов.
// ctxKeyNS — ключ для K8s namespace пользователя в context.Context.
type ctxKeyNS struct{}
// ctxKeyIdentity — ключ для UserIdentity пользователя в context.Context.
type ctxKeyIdentity struct{}
// authTokenFromRequest извлекает Bearer токен из X-Auth-Token или Authorization заголовка.
func authTokenFromRequest(r *http.Request) string {
token := strings.TrimSpace(r.Header.Get("X-Auth-Token"))
if token != "" {
@@ -32,61 +33,57 @@ func authTokenFromRequest(r *http.Request) string {
}
// userNS возвращает namespace пользователя из контекста запроса.
// Устанавливается в authMiddleware после успешной аутентификации.
func (s *Server) userNS(r *http.Request) string {
if ns, ok := r.Context().Value(ctxKeyNS{}).(string); ok && ns != "" {
return ns
}
// Fallback: использовать системный namespace (не должно происходить в prod)
return s.ns
}
// authMiddleware оборачивает handler, добавляя аутентификацию и инициализацию namespace.
// normalizeEnv приводит название стенда к допустимому значению.
func normalizeEnv(env string) string {
switch strings.TrimSpace(strings.ToLower(env)) {
case "prod", "dev", "test":
return strings.TrimSpace(strings.ToLower(env))
}
return "test"
}
// authMiddleware оборачивает handler аутентификацией.
//
// В testMode (FISSION_TEST_MODE=true):
// - Deck API не вызывается
// - X-Test-Sub или X-Auth-Token задают sub → разные namespace-ы для тестирования
//
// В production:
// - X-Auth-Token валидируется через Deck API
// - Namespace вычисляется из JWT claim "sub"
// В testMode: заголовок X-Test-Sub позволяет задать sub напрямую (без токена).
// Иначе: токен передаётся в s.authenticator.Authenticate — детали скрыты за интерфейсом.
func (s *Server) authMiddleware(h http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
var ns string
env := strings.TrimSpace(strings.ToLower(r.Header.Get("X-Auth-Env")))
if _, ok := deckAPIs[env]; !ok {
env = "test"
}
var identity auth.UserIdentity
env := normalizeEnv(r.Header.Get("X-Auth-Env"))
if s.testMode {
sub := strings.TrimSpace(r.Header.Get("X-Test-Sub"))
if sub != "" {
ns = namespaceFromSub(sub)
} else {
token := authTokenFromRequest(r)
resolvedNS, err := s.resolveNamespaceForToken(token, env, true)
if err != nil {
writeJSONError(w, http.StatusUnauthorized, "unauthorized")
if sub := strings.TrimSpace(r.Header.Get("X-Test-Sub")); sub != "" {
identity = auth.UserIdentity{Sub: sub}
ctx := s.contextWithIdentity(r.Context(), identity)
if err := s.nsManager.EnsureUserNS(ctx, auth.NamespaceForSub(sub)); err != nil {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure namespace: %v", err))
return
}
ns = resolvedNS
}
} else {
token := authTokenFromRequest(r)
resolvedNS, err := s.resolveNamespaceForToken(token, env, false)
if err != nil {
writeJSONError(w, http.StatusUnauthorized, "unauthorized")
h(w, r.WithContext(ctx))
return
}
ns = resolvedNS
}
ctx := context.WithValue(r.Context(), ctxKeyNS{}, ns)
token := authTokenFromRequest(r)
var err error
identity, err = s.authenticator.Authenticate(r.Context(), token, env)
if err != nil {
writeJSONError(w, http.StatusUnauthorized, "unauthorized")
return
}
// Гарантируем что namespace + RBAC + quota + netpol существуют.
// EnsureUserNS реализует singleflight + кэш + семафор параллелизма.
if ensureErr := s.nsManager.EnsureUserNS(ctx, ns); ensureErr != nil {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure namespace: %v", ensureErr))
ns := auth.NamespaceForSub(identity.Sub)
ctx := s.contextWithIdentity(r.Context(), identity)
if err := s.nsManager.EnsureUserNS(ctx, ns); err != nil {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure namespace: %v", err))
return
}
@@ -94,72 +91,10 @@ func (s *Server) authMiddleware(h http.HandlerFunc) http.HandlerFunc {
}
}
func namespaceFromSub(sub string) string {
h32 := sha256.Sum256([]byte(sub))
return "fission-" + hex.EncodeToString(h32[:8])
}
func (s *Server) resolveNamespaceForToken(token, env string, allowTestSub bool) (string, error) {
token = strings.TrimSpace(token)
if token == "" {
if allowTestSub {
return "", fmt.Errorf("test mode: X-Test-Sub required")
}
return "", fmt.Errorf("unauthorized")
}
if err := s.validateDeckToken(token, env); err == nil {
ns, nsErr := namespaceFromJWT(token)
if nsErr != nil {
return s.ns, nil
}
return ns, nil
}
if allowTestSub && strings.Contains(token, "@") {
return namespaceFromSub(token), nil
}
return "", fmt.Errorf("invalid token")
}
// namespaceFromJWT декодирует JWT payload (без верификации подписи),
// извлекает claim "sub" и вычисляет namespace: "fission-" + hex(SHA256(sub)[:8]).
//
// Подпись не проверяется — токен уже валидирован через Deck API (validateDeckToken).
// Здесь нам нужен только deterministic namespace name из sub claim.
func namespaceFromJWT(token string) (string, error) {
parts := strings.SplitN(token, ".", 3)
if len(parts) != 3 {
return "", fmt.Errorf("invalid JWT format")
}
payload := parts[1]
// JWT использует base64url без padding — добавляем если нужно
switch len(payload) % 4 {
case 2:
payload += "=="
case 3:
payload += "="
}
// base64url без стандартного padding — пробуем оба варианта
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
// contextWithIdentity кладёт UserIdentity и вычисленный namespace в контекст.
func (s *Server) contextWithIdentity(ctx context.Context, identity auth.UserIdentity) context.Context {
ns := auth.NamespaceForSub(identity.Sub)
ctx = context.WithValue(ctx, ctxKeyNS{}, ns)
ctx = context.WithValue(ctx, ctxKeyIdentity{}, identity)
return ctx
}
-9
View File
@@ -1,9 +0,0 @@
package api
import (
"net/http"
)
func (s *Server) handleAuth(w http.ResponseWriter, r *http.Request) {
// ...existing code...
}
-112
View File
@@ -1,112 +0,0 @@
package api
import (
"net/http"
"net/http/httptest"
"testing"
"time"
)
const validJWTForTests = "eyJhbGciOiJSUzI1NiIsInR5cCI6IkpXVCJ9.eyJzdWIiOiJ0ZXN0LXVzZXItMTIzIn0.signature"
func TestResolveNamespaceForTokenRejectsInvalidTokenOutsideTestMode(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/index.cfm/instances" {
t.Fatalf("unexpected path: %s", r.URL.Path)
}
w.WriteHeader(http.StatusUnauthorized)
}))
defer server.Close()
oldDeckAPIs := deckAPIs
deckAPIs = map[string]string{"test": server.URL}
defer func() { deckAPIs = oldDeckAPIs }()
s := &Server{http: server.Client()}
if _, err := s.resolveNamespaceForToken("user@example.com", "test", false); err == nil {
t.Fatal("expected invalid token to be rejected outside test mode")
}
}
func TestResolveNamespaceForTokenAllowsValidatedJWTInTestMode(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if got := r.Header.Get("Authorization"); got == "" {
t.Fatal("expected Authorization header")
}
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte(`{"results":[]}`))
}))
defer server.Close()
oldDeckAPIs := deckAPIs
deckAPIs = map[string]string{"test": server.URL}
defer func() { deckAPIs = oldDeckAPIs }()
s := &Server{http: server.Client()}
ns, err := s.resolveNamespaceForToken(validJWTForTests, "test", true)
if err != nil {
t.Fatalf("expected JWT to pass in test mode: %v", err)
}
expected, err := namespaceFromJWT(validJWTForTests)
if err != nil {
t.Fatalf("namespaceFromJWT: %v", err)
}
if ns != expected {
t.Fatalf("expected namespace %q, got %q", expected, ns)
}
}
func TestResolveNamespaceForTokenAllowsEmailFallbackOnlyInTestMode(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusUnauthorized)
}))
defer server.Close()
oldDeckAPIs := deckAPIs
deckAPIs = map[string]string{"test": server.URL}
defer func() { deckAPIs = oldDeckAPIs }()
s := &Server{http: server.Client()}
ns, err := s.resolveNamespaceForToken("user@example.com", "test", true)
if err != nil {
t.Fatalf("expected email fallback in test mode: %v", err)
}
if ns != namespaceFromSub("user@example.com") {
t.Fatalf("unexpected namespace: %q", ns)
}
_, err = s.resolveNamespaceForToken("user@example.com", "test", false)
if err == nil {
t.Fatal("expected email fallback to be rejected outside test mode")
}
}
func TestValidateDeckTokenCachesSuccessfulValidation(t *testing.T) {
requests := 0
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
requests++
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte(`{"results":[]}`))
}))
defer server.Close()
oldDeckAPIs := deckAPIs
deckAPIs = map[string]string{"test": server.URL}
defer func() { deckAPIs = oldDeckAPIs }()
s := &Server{http: server.Client()}
if err := s.validateDeckToken(validJWTForTests, "test"); err != nil {
t.Fatalf("first validate failed: %v", err)
}
if err := s.validateDeckToken(validJWTForTests, "test"); err != nil {
t.Fatalf("second validate failed: %v", err)
}
if requests != 1 {
t.Fatalf("expected 1 upstream request because of cache, got %d", requests)
}
if v, ok := s.tokenCache.Load("test:" + validJWTForTests); !ok || time.Now().After(v.(time.Time)) {
t.Fatal("expected token to be cached")
}
}
-19
View File
@@ -1,19 +0,0 @@
package api
import "net/http"
func (s *Server) handleCreateFunction(w http.ResponseWriter, r *http.Request) {
// ...existing code...
}
func (s *Server) handleGetFunction(w http.ResponseWriter, r *http.Request, name string) {
// ...existing code...
}
func (s *Server) handleUpdateFunctionCode(w http.ResponseWriter, r *http.Request, name string) {
// ...existing code...
}
func (s *Server) handleDeleteFunction(w http.ResponseWriter, r *http.Request, name string) {
// ...existing code...
}
+513 -23
View File
@@ -3,6 +3,7 @@ package api
import (
"bytes"
"context"
"encoding/base64"
"encoding/json"
"errors"
"fmt"
@@ -15,6 +16,7 @@ import (
"strings"
"time"
"fission-console/internal/auth"
"fission-console/internal/fission"
"fission-console/internal/model"
"fission-console/internal/runtime"
@@ -36,10 +38,189 @@ const maxCodeSize = 1 << 20
// defaultFunctionInvokeTimeout совпадает с дефолтом Fission для spec.functionTimeout.
const defaultFunctionInvokeTimeout = 60 * time.Second
func buildDeployArchive(lang, code string) ([]byte, error) {
switch lang {
case "nodejs":
return runtime.BuildJSDeployZip(code)
case "php":
return runtime.BuildScriptZip(code, "main.php")
case "ruby":
return runtime.BuildScriptZip(code, "handler.rb")
default:
return []byte(code), nil
}
}
func (s *Server) resolveInvokeTimeout(fn *unstructured.Unstructured) time.Duration {
if fn != nil {
seconds, found, err := unstructured.NestedInt64(fn.Object, "spec", "functionTimeout")
if err == nil && found && seconds > 0 {
return time.Duration(seconds) * time.Second
}
}
if s.invokeTimeout > 0 {
return s.invokeTimeout
}
return defaultFunctionInvokeTimeout
}
func normalizeFunctionTimeout(seconds int64) int64 {
if seconds <= 0 {
return int64(defaultFunctionInvokeTimeout / time.Second)
}
return seconds
}
func normalizeRoute(route string) string {
route = strings.TrimSpace(route)
if route == "" || route == "/" {
return "/"
}
if !strings.HasPrefix(route, "/") {
return "/" + route
}
return route
}
func routeAllowsMethod(methods []string, method string) bool {
if len(methods) == 0 {
return true
}
for _, candidate := range methods {
if strings.EqualFold(strings.TrimSpace(candidate), method) {
return true
}
}
return false
}
func appendUniqueMethods(dst []string, src []string) []string {
for _, method := range src {
method = strings.ToUpper(strings.TrimSpace(method))
if method == "" {
continue
}
seen := false
for _, existing := range dst {
if existing == method {
seen = true
break
}
}
if !seen {
dst = append(dst, method)
}
}
return dst
}
func shouldForwardRequestBody(method string) bool {
switch method {
case http.MethodGet, http.MethodHead:
return false
default:
return true
}
}
func copyProxyRequestHeaders(dst, src http.Header) {
for key, values := range src {
switch http.CanonicalHeaderKey(key) {
case "Authorization", "X-Auth-Token", "X-Auth-Env", "Host", "Content-Length":
continue
}
for _, value := range values {
dst.Add(key, value)
}
}
}
func copyProxyResponseHeaders(dst, src http.Header) {
for key, values := range src {
for _, value := range values {
dst.Add(key, value)
}
}
}
func doRequestWithContextTimeout(client *http.Client, req *http.Request) (*http.Response, error) {
if client == nil {
return http.DefaultClient.Do(req)
}
invokeClient := *client
// Для invoke реальный лимит должен определяться context timeout функции,
// а не общим HTTP timeout console.
invokeClient.Timeout = 0
return invokeClient.Do(req)
}
// handleFunctionsRoot обрабатывает запросы к /console/api/functions без имени функции.
// GET → список всех функций, POST → создать новую.
func (s *Server) handleFunctionsRoot(w http.ResponseWriter, r *http.Request) {
switch r.Method {
case http.MethodGet:
s.handleList(fission.FunctionGVR)(w, r)
case http.MethodPost:
s.handleCreateFunction(w, r)
default:
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
}
}
// handleFunctionsAction обрабатывает запросы к /console/api/functions/:name[/action].
// Парсит имя функции и опциональный sub-path ("code", "invoke").
func (s *Server) handleFunctionsAction(w http.ResponseWriter, r *http.Request) {
// Убираем оба возможных префикса (legacy /api/ и основной /console/api/)
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 {
// /functions/:name — CRUD операции с конкретной функцией
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
}
// /functions/:name/code — обновление кода
if len(parts) == 2 && parts[1] == "code" && r.Method == http.MethodPut {
s.handleUpdateFunctionCode(w, r, name)
return
}
// /functions/:name/invoke — вызов функции
if len(parts) == 2 && parts[1] == "invoke" && r.Method == http.MethodPost {
s.handleInvokeFunction(w, r, name)
return
}
http.NotFound(w, r)
}
// handleCreateFunction создаёт новую функцию: Package + Function + HTTPTrigger.
//
// Порядок создания: Package → Function → HTTPTrigger.
// При ошибке на любом шаге откатываем уже созданные объекты (best-effort).
// TTL парсится ДО создания объектов — невалидный TTL не оставляет мусор.
func (s *Server) handleCreateFunction(w http.ResponseWriter, r *http.Request) {
ns := s.userNS(r)
@@ -121,19 +302,140 @@ func (s *Server) handleCreateFunction(w http.ResponseWriter, r *http.Request) {
return
}
result, err := fission.CreateFunction(ctx, s.dyn, ns, req)
if err != nil {
status := http.StatusBadGateway
if apierrors.IsAlreadyExists(err) {
status = http.StatusConflict
} else if apierrors.IsInvalid(err) {
status = http.StatusBadRequest
pkgName := req.Name + "-pkg"
triggerName := req.Name + "-route"
methodValues := make([]any, 0, len(req.Methods))
for _, method := range req.Methods {
methodValues = append(methodValues, method)
}
// Строим Package spec в зависимости от языка:
// - Go: source package → builder job компилирует в .so плагин
// - Node.js: deployment zip с ESM wrapper (package.json + main.js)
// - Остальные: deployment literal с кодом напрямую
var pkgSpec map[string]any
if req.Language == "go" {
srcZip, err := runtime.BuildGoSourceZip(req.Code)
if err != nil {
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("build go source archive: %v", err))
return
}
writeJSONError(w, status, err.Error())
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 {
deployBytes, archiveErr := buildDeployArchive(req.Language, req.Code)
if archiveErr != nil {
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build %s archive: %v", req.Language, archiveErr))
return
}
pkgSpec = map[string]any{
"deployment": map[string]any{"type": "literal", "literal": base64.StdEncoding.EncodeToString(deployBytes)},
"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{
"fission-console/language": req.Language,
}
now := time.Now().UTC()
fnAnnotations[functionCreatedAtAnnotation] = now.Format(time.RFC3339)
fnAnnotations[functionUpdatedAtAnnotation] = now.Format(time.RFC3339)
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(fission.PackageGVR).Namespace(ns).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
}
writeAnyJSON(w, http.StatusCreated, result)
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},
"functionTimeout": req.Timeout,
"InvokeStrategy": map[string]any{
"ExecutionStrategy": map[string]any{"ExecutorType": "poolmgr"},
"StrategyType": "execution",
},
"package": map[string]any{
"packageref": map[string]any{"name": pkgName, "namespace": ns},
"functionName": req.Entrypoint,
},
},
}}
if _, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Create(ctx, fn, metav1.CreateOptions{}); err != nil {
// Откатываем Package если Function не создалась
_ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).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(fission.HTTPTrigGVR).Namespace(ns).Create(ctx, httpTrigger, metav1.CreateOptions{}); err != nil {
// Откатываем Function и Package
_ = s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Delete(ctx, req.Name, metav1.DeleteOptions{})
_ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).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"],
})
}
// handleGetFunction возвращает детали функции: код, environment, route, methods.
@@ -165,7 +467,7 @@ func (s *Server) handleGetFunction(w http.ResponseWriter, r *http.Request, name
if packageName != "" {
pkg, pkgErr := s.dyn.Resource(fission.PackageGVR).Namespace(ns).Get(ctx, packageName, metav1.GetOptions{})
if pkgErr == nil {
code = extractPackageSourceCode(ctx, s.http, pkg)
code = extractPackageSourceCode(ctx, s, pkg)
}
}
@@ -201,6 +503,10 @@ func (s *Server) handleGetFunction(w http.ResponseWriter, r *http.Request, name
})
}
// handleUpdateFunctionCode обновляет код уже существующей функции.
// Создаёт НОВЫЙ Package (вместо обновления старого) чтобы executor сбросил кэш:
// executor кэширует function service по functionUid и не видит изменений в том же Package.
// Новое имя пакета гарантирует cache miss в executor.
func (s *Server) handleUpdateFunctionCode(w http.ResponseWriter, r *http.Request, name string) {
var req model.UpdateCodeRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
@@ -212,25 +518,109 @@ func (s *Server) handleUpdateFunctionCode(w http.ResponseWriter, r *http.Request
writeJSONError(w, http.StatusBadRequest, "code is required")
return
}
if req.Timeout <= 0 {
req.Timeout = int64(defaultFunctionInvokeTimeout / time.Second)
}
ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second)
defer cancel()
ns := s.userNS(r)
result, err := fission.UpdateFunctionCode(ctx, s.dyn, ns, name, req)
fn, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Get(ctx, name, metav1.GetOptions{})
if err != nil {
status := http.StatusBadGateway
if apierrors.IsNotFound(err) {
status = http.StatusNotFound
}
writeJSONError(w, status, err.Error())
writeJSONError(w, status, fmt.Sprintf("get function %q: %v", name, err))
return
}
writeAnyJSON(w, http.StatusOK, result)
oldPkgName, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name")
// Определяем язык из аннотации — нужен для правильной упаковки
lang, _, _ := unstructured.NestedString(fn.Object, "metadata", "annotations", "fission-console/language")
deployBytes, archiveErr := buildDeployArchive(lang, req.Code)
if archiveErr != nil {
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build %s archive: %v", lang, archiveErr))
return
}
// Создаём новый Package с уникальным именем.
// Это единственный способ сбросить кэш executor: он кэширует по functionUid и
// не замечает изменений в существующем Package.
envName, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name")
createdAt := func() time.Time {
ann := fn.GetAnnotations()
if ann != nil {
if v := strings.TrimSpace(ann[functionCreatedAtAnnotation]); v != "" {
if ts, err := parseRFC3339(v); err == nil {
return ts.UTC()
}
}
}
if ts := fn.GetCreationTimestamp(); !ts.IsZero() {
return ts.UTC()
}
return time.Time{}
}()
now := time.Now().UTC()
newPkgName := name + "-pkg-" + strconv.FormatInt(time.Now().UnixMilli(), 36)
newPkg := &unstructured.Unstructured{Object: map[string]any{
"apiVersion": "fission.io/v1",
"kind": "Package",
"metadata": map[string]any{"name": newPkgName, "namespace": ns},
"spec": map[string]any{
"deployment": map[string]any{"type": "literal", "literal": base64.StdEncoding.EncodeToString(deployBytes)},
"environment": map[string]any{"name": envName, "namespace": ns},
"source": map[string]any{},
},
}}
createdPkg, err := s.dyn.Resource(fission.PackageGVR).Namespace(ns).Create(ctx, newPkg, metav1.CreateOptions{})
if err != nil {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create new package: %v", err))
return
}
// Обновляем Function на новый Package
if err := unstructured.SetNestedField(fn.Object, normalizeFunctionTimeout(req.Timeout), "spec", "functionTimeout"); err != nil {
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set function timeout: %v", err))
_ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, newPkgName, metav1.DeleteOptions{})
return
}
ensureFunctionTimestamps(fn, now)
if createdAt.IsZero() {
createdAt = now
}
fnAnnotations := fn.GetAnnotations()
if fnAnnotations == nil {
fnAnnotations = map[string]string{}
}
fnAnnotations[functionCreatedAtAnnotation] = createdAt.UTC().Format(time.RFC3339)
fnAnnotations[functionUpdatedAtAnnotation] = now.Format(time.RFC3339)
fn.SetAnnotations(fnAnnotations)
if err := unstructured.SetNestedField(fn.Object, map[string]any{
"name": newPkgName,
"namespace": ns,
"resourceversion": createdPkg.GetResourceVersion(),
}, "spec", "package", "packageref"); err != nil {
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set function packageref: %v", err))
_ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, newPkgName, metav1.DeleteOptions{})
return
}
if _, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Update(ctx, fn, metav1.UpdateOptions{}); err != nil {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("update function %q: %v", name, err))
_ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, newPkgName, metav1.DeleteOptions{})
return
}
// Удаляем старый Package (best effort)
if oldPkgName != "" && oldPkgName != newPkgName {
_ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, oldPkgName, metav1.DeleteOptions{})
}
writeAnyJSON(w, http.StatusOK, map[string]any{
"updated": true,
"package": newPkgName,
})
}
// handleInvokeFunction вызывает функцию через Fission router.
@@ -562,26 +952,108 @@ func (s *Server) handleInvokeRoute(w http.ResponseWriter, r *http.Request) {
_, _ = io.Copy(w, resp.Body)
}
// handleDeleteFunction удаляет функцию и связанные объекты: HTTPTrigger, TimeTrigger, Package.
// После удаления вызывает CleanupEnvironmentIfUnused — убирает environment если язык больше не используется.
func (s *Server) handleDeleteFunction(w http.ResponseWriter, r *http.Request, name string) {
ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second)
defer cancel()
ns := s.userNS(r)
if err := fission.DeleteFunction(ctx, s.dyn, ns, name); err != nil {
status := http.StatusBadGateway
// Получаем Function чтобы знать pkgName и envName для cleanup
var pkgName, envName string
fn, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Get(ctx, name, metav1.GetOptions{})
if err != nil {
if apierrors.IsNotFound(err) {
status = http.StatusNotFound
writeJSONError(w, http.StatusNotFound, fmt.Sprintf("function %q not found", name))
return
}
writeJSONError(w, status, err.Error())
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")
// Удаляем связанные HTTPTrigger-ы
triggers, err := s.dyn.Resource(fission.HTTPTrigGVR).Namespace(ns).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(fission.HTTPTrigGVR).Namespace(ns).Delete(ctx, trig.GetName(), metav1.DeleteOptions{})
}
}
}
// Удаляем связанные TimeTrigger-ы
if triggers, err := s.dyn.Resource(fission.TimeTrigGVR).Namespace(ns).List(ctx, metav1.ListOptions{}); err == nil {
for _, trig := range triggers.Items {
refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name")
if refName == name {
_ = s.dyn.Resource(fission.TimeTrigGVR).Namespace(ns).Delete(ctx, trig.GetName(), metav1.DeleteOptions{})
}
}
}
if err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Delete(ctx, name, metav1.DeleteOptions{}); err != nil && !apierrors.IsNotFound(err) {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("delete function %q: %v", name, err))
return
}
writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name})
if pkgName != "" {
if err := s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, pkgName, metav1.DeleteOptions{}); err != nil && !apierrors.IsNotFound(err) {
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("delete package %q: %v", pkgName, err))
return
}
}
// Убираем environment pool pods если язык больше не используется (best-effort)
if envName != "" {
cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cleanupCancel()
fission.CleanupEnvironmentIfUnused(cleanupCtx, s.dyn, ns, envName)
}
// (reconciler NS удалён — за FISSION_RESOURCE_NAMESPACES теперь отвечает Layer 1 NSWatcher)
writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name, "package": pkgName})
}
// handleAuth обрабатывает POST /console/api/auth.
// Валидирует токен, создаёт namespace, возвращает namespace пользователя.
// Валидирует токен через authenticator, создаёт namespace, возвращает namespace + email.
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 := normalizeEnv(body.Env)
identity, err := s.authenticator.Authenticate(r.Context(), body.Token, env)
if err != nil {
writeJSONError(w, http.StatusUnauthorized, "invalid token")
return
}
ns := auth.NamespaceForSub(identity.Sub)
ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second)
defer cancel()
if ensureErr := s.nsManager.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, "email": identity.Email})
}
// parseTTL парсит строку TTL и возвращает время истечения.
// Поддерживаемые форматы: Go duration (1h, 30m, 24h) и дни (1d, 7d, 30d).
@@ -606,4 +1078,22 @@ func parseTTL(ttl string) (time.Time, error) {
// normalizeMethods приводит список HTTP методов к верхнему регистру, убирает дубли.
// Если список пустой или все элементы пустые — возвращает ["GET"].
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
}
-48
View File
@@ -1,48 +0,0 @@
package api
import (
"net/http"
"time"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
)
func (s *Server) handleInvokeFunction(w http.ResponseWriter, r *http.Request, name string) {
// ...existing code...
}
func (s *Server) handleFissionFunctionGateway(w http.ResponseWriter, r *http.Request) {
// ...existing code...
}
func (s *Server) invokeInternalFunction(w http.ResponseWriter, r *http.Request, namespace, functionName, extraPath string) {
// ...existing code...
}
func (s *Server) handleInvokeRoute(w http.ResponseWriter, r *http.Request) {
// ...existing code...
}
func buildInternalInvokeURL(routerURL, namespace, functionName string) string {
// ...existing code...
}
func resolveInvokeTimeout(fn *unstructured.Unstructured) time.Duration {
// ...existing code...
}
func shouldForwardRequestBody(method string) bool {
// ...existing code...
}
func copyProxyRequestHeaders(dst, src http.Header) {
// ...existing code...
}
func copyProxyResponseHeaders(dst, src http.Header) {
// ...existing code...
}
func doRequestWithContextTimeout(client *http.Client, req *http.Request) (*http.Response, error) {
// ...existing code...
}
+4 -7
View File
@@ -24,7 +24,7 @@ import (
//
// Декодирование: если данные — валидный UTF-8, возвращаем как есть.
// Если это zip — ищем известные файлы (main.py, handler.rb и т.д.).
func extractPackageSourceCode(ctx context.Context, httpClient *http.Client, pkg *unstructured.Unstructured) string {
func extractPackageSourceCode(ctx context.Context, s *Server, pkg *unstructured.Unstructured) string {
literalPaths := [][]string{
{"spec", "source", "literal"},
{"spec", "deployment", "literal"},
@@ -49,7 +49,7 @@ func extractPackageSourceCode(ctx context.Context, httpClient *http.Client, pkg
if !found || strings.TrimSpace(urlValue) == "" {
continue
}
archiveBytes, err := fetchPackageArchive(ctx, httpClient, urlValue)
archiveBytes, err := fetchPackageArchive(ctx, s, urlValue)
if err != nil {
continue
}
@@ -62,15 +62,12 @@ func extractPackageSourceCode(ctx context.Context, httpClient *http.Client, pkg
}
// fetchPackageArchive скачивает архив функции по URL из Fission storage.
func fetchPackageArchive(ctx context.Context, httpClient *http.Client, archiveURL string) ([]byte, error) {
func fetchPackageArchive(ctx context.Context, s *Server, archiveURL string) ([]byte, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, archiveURL, nil)
if err != nil {
return nil, err
}
if httpClient == nil {
httpClient = http.DefaultClient
}
resp, err := httpClient.Do(req)
resp, err := s.http.Do(req)
if err != nil {
return nil, err
}
+7 -54
View File
@@ -5,13 +5,13 @@ import (
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"os"
"strings"
"sync"
"time"
"fission-console/internal/auth"
"fission-console/internal/cloud"
"fission-console/internal/fission"
"fission-console/ui"
@@ -22,20 +22,12 @@ import (
"k8s.io/client-go/dynamic"
)
// deckAPIs — карта окружений Deck API.
// Ключ используется в X-Auth-Env заголовке для выбора нужного сервера.
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",
}
// defaultSATokenPath — путь к service account токену внутри pod-а.
// Используется для авторизации запросов от console к Fission router.
const defaultSATokenPath = "/var/run/secrets/kubernetes.io/serviceaccount/token"
// Server — основная структура HTTP сервера.
// Содержит все зависимости: kubernetes client, конфиги, кэши токенов.
// Содержит все зависимости: kubernetes client, конфиги, кэши.
type Server struct {
dyn dynamic.Interface
ns string // системный namespace (fallback, обычно "fission")
@@ -44,7 +36,7 @@ type Server struct {
saTokenPath string
invokeTimeout time.Duration
testMode bool // FISSION_TEST_MODE=true — пропускает deck auth
testMode bool // FISSION_TEST_MODE=true — разрешает X-Test-Sub shortcut
authUser string
authPass string
@@ -59,9 +51,8 @@ type Server struct {
cachedJWT string
tokenExpAt time.Time
// tokenCache кэширует результаты валидации Deck токенов.
// Ключ: "env:token", значение: time.Time — когда кэш истекает.
tokenCache sync.Map
// authenticator — слой аутентификации. Сервер не знает деталей реализации.
authenticator auth.Authenticator
// nsManager управляет жизненным циклом пользовательских namespace-ов.
nsManager *cloud.NSManager
@@ -78,6 +69,7 @@ type Config struct {
AuthUser string
AuthPass string
TestMode bool
Authenticator auth.Authenticator // слой аутентификации
LLMUrl string
LLMKey string
}
@@ -94,6 +86,7 @@ func NewServer(cfg Config) *Server {
authUser: cfg.AuthUser,
authPass: cfg.AuthPass,
testMode: cfg.TestMode,
authenticator: cfg.Authenticator,
llmURL: cfg.LLMUrl,
llmKey: cfg.LLMKey,
nsManager: cloud.NewNSManager(cfg.Dyn),
@@ -274,43 +267,3 @@ func (s *Server) getRouterToken() string {
s.tokenExpAt = time.Now().Add(100 * time.Second)
return s.cachedJWT
}
// validateDeckToken проверяет токен через Deck API с кэшированием результата на 5 минут.
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 deck 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
}
-25
View File
@@ -1,25 +0,0 @@
package api
import (
"strings"
)
func normalizeFunctionTimeout(seconds int64) int64 {
// ...existing code...
}
func normalizeRoute(route string) string {
// ...existing code...
}
func normalizeMethods(in []string) []string {
// ...existing code...
}
func routeAllowsMethod(methods []string, method string) bool {
// ...existing code...
}
func appendUniqueMethods(dst []string, src []string) []string {
// ...existing code...
}