diff --git a/console/main.go b/console/main.go index d60fd97..8ae223b 100644 --- a/console/main.go +++ b/console/main.go @@ -4,7 +4,9 @@ import ( "archive/zip" "bytes" "context" + "crypto/sha256" "encoding/base64" + "encoding/hex" "encoding/json" "errors" "fmt" @@ -36,6 +38,7 @@ var ( 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"} ) const defaultSATokenPath = "/var/run/secrets/kubernetes.io/serviceaccount/token" @@ -160,7 +163,13 @@ func main() { writeJSONError(w, http.StatusUnauthorized, "unauthorized") return } - h(w, r) + ns, err := namespaceFromJWT(token) + if err != nil { + log.Printf("namespaceFromJWT: %v", err) + ns = s.ns + } + ctx := context.WithValue(r.Context(), ctxKeyNS{}, ns) + h(w, r.WithContext(ctx)) } } @@ -238,6 +247,7 @@ func (s *server) handleFunctionsAction(w http.ResponseWriter, r *http.Request) { 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 @@ -260,10 +270,10 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { envName := "console-" + req.Language + "-env" ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) defer cancel() - _, err := s.dyn.Resource(environmentGVR).Namespace(s.ns).Get(ctx, envName, metav1.GetOptions{}) + _, err := s.dyn.Resource(environmentGVR).Namespace(s.userNS(r)).Get(ctx, envName, metav1.GetOptions{}) if apierrors.IsNotFound(err) { - env := s.buildLangEnvironment(envName, langDef) - if _, err := s.dyn.Resource(environmentGVR).Namespace(s.ns).Create(ctx, env, metav1.CreateOptions{}); err != nil { + env := s.buildLangEnvironment(envName, s.userNS(r), langDef) + if _, err := s.dyn.Resource(environmentGVR).Namespace(s.userNS(r)).Create(ctx, env, metav1.CreateOptions{}); err != nil { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create environment %q: %v", envName, err)) return } @@ -292,7 +302,7 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second) defer cancel() - if _, err := s.dyn.Resource(environmentGVR).Namespace(s.ns).Get(ctx, req.Environment, metav1.GetOptions{}); err != nil { + 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 } @@ -321,7 +331,7 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { "deployment": map[string]any{}, "environment": map[string]any{ "name": req.Environment, - "namespace": s.ns, + "namespace": ns, }, "buildcommand": "build", } @@ -334,7 +344,7 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { }, "environment": map[string]any{ "name": req.Environment, - "namespace": s.ns, + "namespace": ns, }, "source": map[string]any{}, } @@ -345,12 +355,12 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { "kind": "Package", "metadata": map[string]any{ "name": pkgName, - "namespace": s.ns, + "namespace": ns, }, "spec": pkgSpec, }} - if _, err := s.dyn.Resource(packageGVR).Namespace(s.ns).Create(ctx, pkg, metav1.CreateOptions{}); err != nil { + if _, err := s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Create(ctx, pkg, metav1.CreateOptions{}); err != nil { writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create package: %v", err)) return } @@ -360,26 +370,26 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { "kind": "Function", "metadata": map[string]any{ "name": req.Name, - "namespace": s.ns, + "namespace": ns, }, "spec": map[string]any{ "environment": map[string]any{ "name": req.Environment, - "namespace": s.ns, + "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.ns}, + "packageref": map[string]any{"name": pkgName, "namespace": s.userNS(r)}, "functionName": req.Entrypoint, }, }, }} - if _, err := s.dyn.Resource(functionGVR).Namespace(s.ns).Create(ctx, fn, metav1.CreateOptions{}); err != nil { - _ = s.dyn.Resource(packageGVR).Namespace(s.ns).Delete(ctx, pkgName, metav1.DeleteOptions{}) + 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{}) writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create function: %v", err)) return } @@ -389,7 +399,7 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { "kind": "HTTPTrigger", "metadata": map[string]any{ "name": triggerName, - "namespace": s.ns, + "namespace": ns, }, "spec": map[string]any{ "relativeurl": req.Route, @@ -402,9 +412,9 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { }, }} - if _, err := s.dyn.Resource(httpTrigGVR).Namespace(s.ns).Create(ctx, httpTrigger, metav1.CreateOptions{}); err != nil { - _ = s.dyn.Resource(functionGVR).Namespace(s.ns).Delete(ctx, req.Name, metav1.DeleteOptions{}) - _ = s.dyn.Resource(packageGVR).Namespace(s.ns).Delete(ctx, pkgName, metav1.DeleteOptions{}) + 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 } @@ -417,7 +427,7 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) { }) } -func (s *server) buildLangEnvironment(name string, def langEnvDef) *unstructured.Unstructured { +func (s *server) buildLangEnvironment(name string, ns string, def langEnvDef) *unstructured.Unstructured { spec := map[string]any{ "version": int64(3), "runtime": map[string]any{ @@ -436,7 +446,7 @@ func (s *server) buildLangEnvironment(name string, def langEnvDef) *unstructured "kind": "Environment", "metadata": map[string]any{ "name": name, - "namespace": s.ns, + "namespace": ns, }, "spec": spec, }} @@ -473,7 +483,7 @@ func (s *server) handleGetFunction(w http.ResponseWriter, r *http.Request, name ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) defer cancel() - fn, err := s.dyn.Resource(functionGVR).Namespace(s.ns).Get(ctx, name, metav1.GetOptions{}) + 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) { @@ -489,7 +499,7 @@ func (s *server) handleGetFunction(w http.ResponseWriter, r *http.Request, name code := "" if packageName != "" { - pkg, pkgErr := s.dyn.Resource(packageGVR).Namespace(s.ns).Get(ctx, packageName, metav1.GetOptions{}) + pkg, pkgErr := s.dyn.Resource(packageGVR).Namespace(s.userNS(r)).Get(ctx, packageName, metav1.GetOptions{}) if pkgErr == nil { code = s.extractPackageSourceCode(ctx, pkg) } @@ -497,7 +507,7 @@ func (s *server) handleGetFunction(w http.ResponseWriter, r *http.Request, name route := "" methods := []string{} - triggers, trigErr := s.dyn.Resource(httpTrigGVR).Namespace(s.ns).List(ctx, metav1.ListOptions{}) + 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") @@ -512,7 +522,7 @@ func (s *server) handleGetFunction(w http.ResponseWriter, r *http.Request, name writeAnyJSON(w, http.StatusOK, map[string]any{ "name": name, - "namespace": s.ns, + "namespace": s.userNS(r), "environment": environment, "package": packageName, "entrypoint": entrypoint, @@ -538,7 +548,7 @@ func (s *server) handleUpdateFunctionCode(w http.ResponseWriter, r *http.Request ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second) defer cancel() - fn, err := s.dyn.Resource(functionGVR).Namespace(s.ns).Get(ctx, name, metav1.GetOptions{}) + 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) { @@ -554,7 +564,7 @@ func (s *server) handleUpdateFunctionCode(w http.ResponseWriter, r *http.Request return } - pkg, err := s.dyn.Resource(packageGVR).Namespace(s.ns).Get(ctx, pkgName, metav1.GetOptions{}) + 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) { @@ -570,12 +580,12 @@ func (s *server) handleUpdateFunctionCode(w http.ResponseWriter, r *http.Request return } - if _, err := s.dyn.Resource(packageGVR).Namespace(s.ns).Update(ctx, pkg, metav1.UpdateOptions{}); err != nil { + 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.ns).Get(ctx, pkgName, metav1.GetOptions{}) + 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 @@ -586,7 +596,7 @@ func (s *server) handleUpdateFunctionCode(w http.ResponseWriter, r *http.Request return } - if _, err := s.dyn.Resource(functionGVR).Namespace(s.ns).Update(ctx, fn, metav1.UpdateOptions{}); err != nil { + 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 } @@ -614,7 +624,7 @@ func (s *server) handleInvokeFunction(w http.ResponseWriter, r *http.Request, na invokeURL := fmt.Sprintf("%s/fission-function/v2/functions/%s", s.routerURL, name) invokeMethod := http.MethodPost - triggers, err := s.dyn.Resource(httpTrigGVR).Namespace(s.ns).List(ctx, metav1.ListOptions{}) + 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") @@ -741,6 +751,90 @@ func (s *server) getRouterToken() string { 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 если не существуют. +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{}) + if err != nil && !apierrors.IsAlreadyExists(err) { + return fmt.Errorf("create namespace %s: %w", ns, err) + } + + // 2. Создать shared environments для всех языков + for lang, def := range langEnvMap { + envName := lang + "-env" + _, getErr := s.dyn.Resource(environmentGVR).Namespace(ns).Get(ctx, envName, metav1.GetOptions{}) + if getErr == nil { + continue // уже существует + } + if !apierrors.IsNotFound(getErr) { + continue // другая ошибка — пропускаем + } + env := s.buildLangEnvironment(envName, ns, def) + if _, createErr := s.dyn.Resource(environmentGVR).Namespace(ns).Create(ctx, env, metav1.CreateOptions{}); createErr != nil && !apierrors.IsAlreadyExists(createErr) { + log.Printf("ensureUserNamespace: create env %s/%s: %v", ns, envName, createErr) + } + } + return nil +} + func (s *server) validateDeckToken(token, env string) error { cacheKey := env + ":" + token if v, ok := s.tokenCache.Load(cacheKey); ok { @@ -794,8 +888,18 @@ func (s *server) handleAuth(w http.ResponseWriter, r *http.Request) { writeJSONError(w, http.StatusUnauthorized, "invalid token") return } + ns, err := namespaceFromJWT(body.Token) + if err != nil { + log.Printf("handleAuth: namespaceFromJWT: %v", err) + ns = s.ns + } + ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) + defer cancel() + if ensureErr := s.ensureUserNamespace(ctx, ns); ensureErr != nil { + log.Printf("handleAuth: ensureUserNamespace %s: %v", ns, ensureErr) + } w.Header().Set("Content-Type", "application/json; charset=utf-8") - _ = json.NewEncoder(w).Encode(map[string]any{"ok": true, "env": env}) + _ = 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) { @@ -803,28 +907,28 @@ func (s *server) handleDeleteFunction(w http.ResponseWriter, r *http.Request, na defer cancel() var pkgName string - fn, err := s.dyn.Resource(functionGVR).Namespace(s.ns).Get(ctx, name, metav1.GetOptions{}) + fn, err := s.dyn.Resource(functionGVR).Namespace(s.userNS(r)).Get(ctx, name, metav1.GetOptions{}) if err == nil { pkgName, _, _ = unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name") } - triggers, err := s.dyn.Resource(httpTrigGVR).Namespace(s.ns).List(ctx, metav1.ListOptions{}) + 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.ns).Delete(ctx, trig.GetName(), metav1.DeleteOptions{}) + _ = s.dyn.Resource(httpTrigGVR).Namespace(s.userNS(r)).Delete(ctx, trig.GetName(), metav1.DeleteOptions{}) } } } - if err := s.dyn.Resource(functionGVR).Namespace(s.ns).Delete(ctx, name, metav1.DeleteOptions{}); err != nil && !apierrors.IsNotFound(err) { + 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.ns).Delete(ctx, pkgName, metav1.DeleteOptions{}); err != nil && !apierrors.IsNotFound(err) { + 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 } @@ -862,7 +966,7 @@ func (s *server) handleList(gvr schema.GroupVersionResource) http.HandlerFunc { ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) defer cancel() - list, err := s.dyn.Resource(gvr).Namespace(s.ns).List(ctx, metav1.ListOptions{}) + 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