feat: namespace isolation per user JWT sub claim

- Auth middleware extracts user namespace from JWT sub (SHA256 → hex8)
- All K8s operations use per-user fission-{hex} namespace
- handleAuth calls ensureUserNamespace to create NS + shared envs
- handleAuth returns namespace in response
- buildLangEnvironment accepts ns string parameter
- Fallback to server default NS if JWT parse fails
This commit is contained in:
Naeel
2026-04-19 09:54:57 +03:00
parent 50ba09e1b6
commit 8aad0cec58
+141 -37
View File
@@ -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