fix: semaphore+singleflight for ensureUserNS, fix 504 on 10 parallel new users (v0.8.12)
This commit is contained in:
@@ -46,7 +46,7 @@ spec:
|
||||
serviceAccountName: fission-console
|
||||
containers:
|
||||
- name: console
|
||||
image: naeel/fission-console:v0.3.4
|
||||
image: naeel/fission-console:v0.8.12
|
||||
ports:
|
||||
- containerPort: 8090
|
||||
env:
|
||||
|
||||
+387
-62
@@ -56,6 +56,13 @@ var deckAPIs = map[string]string{
|
||||
"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
|
||||
@@ -77,6 +84,22 @@ type server struct {
|
||||
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 {
|
||||
@@ -145,6 +168,10 @@ func main() {
|
||||
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()
|
||||
@@ -208,6 +235,14 @@ func main() {
|
||||
}
|
||||
|
||||
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))
|
||||
}
|
||||
}
|
||||
@@ -232,6 +267,7 @@ func main() {
|
||||
|
||||
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())
|
||||
}
|
||||
|
||||
@@ -305,10 +341,10 @@ func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
// Гарантируем что namespace + RBAC созданы до любых операций с ресурсами.
|
||||
// handleAuth делает это при логине, но в test mode или при прямом вызове API
|
||||
// namespace может отсутствовать — создаём idempotent.
|
||||
nsCtx, nsCancel := context.WithTimeout(r.Context(), 30*time.Second)
|
||||
// namespace может отсутствовать — создаём idempotent через кэш+singleflight+семафор.
|
||||
nsCtx, nsCancel := context.WithTimeout(r.Context(), 60*time.Second)
|
||||
defer nsCancel()
|
||||
if err := s.ensureUserNamespace(nsCtx, ns); err != nil {
|
||||
if err := s.ensureUserNS(nsCtx, ns); err != nil {
|
||||
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("ensure namespace: %v", err))
|
||||
return
|
||||
}
|
||||
@@ -725,6 +761,12 @@ func (s *server) reapExpiredFunctionsInNS(ctx context.Context, ns string, now ti
|
||||
}
|
||||
}
|
||||
|
||||
// Сигналим reconciler — он проверит все namespace'ы и почистит список.
|
||||
select {
|
||||
case s.nsReconcileCh <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
|
||||
// Сканируем orphan packages — пакеты без соответствующей функции
|
||||
// (могут остаться если под упал в середине удаления)
|
||||
packages, pkgListErr := s.dyn.Resource(packageGVR).Namespace(ns).List(ctx, metav1.ListOptions{})
|
||||
@@ -1194,45 +1236,101 @@ func namespaceFromJWT(token string) (string, error) {
|
||||
}
|
||||
|
||||
// ensureUserNamespace создаёт K8s namespace и shared environments если не существуют.
|
||||
// addNSToFission динамически добавляет namespace в FISSION_RESOURCE_NAMESPACES у всех Fission deployments.
|
||||
// После создания сигналит nsReconciler который синхронизирует FISSION_RESOURCE_NAMESPACES.
|
||||
//
|
||||
// Зачем это нужно: Fission компоненты смотрят только те namespace-ы, что указаны в
|
||||
// FISSION_RESOURCE_NAMESPACES. Если не добавить новый namespace — executor/router не будут
|
||||
// создавать пулы и обрабатывать триггеры, invoke вернёт 404.
|
||||
// Зачем 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. Читаем текущее значение FISSION_RESOURCE_NAMESPACES из router deployment.
|
||||
// 2. Если namespace уже в списке — выходим (idempotent).
|
||||
// 3. Иначе добавляем namespace к списку и патчим все Fission deployments через StrategicMergePatch.
|
||||
//
|
||||
// StrategicMergePatch позволяет обновить только одну env переменную не трогая остальные.
|
||||
func (s *server) addNSToFission(ctx context.Context, ns string) error {
|
||||
// 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"
|
||||
}
|
||||
fissionDeployments := []string{"router", "executor", "buildermgr", "kubewatcher", "timer"}
|
||||
|
||||
// Читаем текущее значение FISSION_RESOURCE_NAMESPACES из router deployment.
|
||||
// Используем router как источник истины — он первым получает изменения.
|
||||
routerDep, err := s.dyn.Resource(deploymentGVR).Namespace(fissionNS).Get(ctx, "router", metav1.GetOptions{})
|
||||
// Шаг 1: реальные namespace'ы с нашей меткой
|
||||
nsList, err := s.dyn.Resource(namespaceGVR).List(ctx, metav1.ListOptions{
|
||||
LabelSelector: "managed-by=fission-console",
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("get router deployment: %w", err)
|
||||
log.Printf("nsReconciler: list namespaces: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
currentVal := "default" // fallback если переменная не найдена
|
||||
containerName := "router" // имя контейнера нужно для StrategicMergePatch
|
||||
// Шаг 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
|
||||
}
|
||||
// Запоминаем реальное имя контейнера — оно используется как merge key в StrategicMergePatch.
|
||||
// Без точного имени патч создаст дублирующий контейнер вместо обновления существующего.
|
||||
if n, ok := cont["name"].(string); ok {
|
||||
containerName = n
|
||||
}
|
||||
envs, _, _ := unstructured.NestedSlice(cont, "env")
|
||||
for _, e := range envs {
|
||||
env, ok := e.(map[string]any)
|
||||
@@ -1245,27 +1343,47 @@ func (s *server) addNSToFission(ctx context.Context, ns string) error {
|
||||
}
|
||||
}
|
||||
}
|
||||
break // берём только первый контейнер
|
||||
break
|
||||
}
|
||||
|
||||
// Проверяем что namespace ещё не в списке
|
||||
for _, existing := range strings.Split(currentVal, ",") {
|
||||
if strings.TrimSpace(existing) == ns {
|
||||
return nil // уже есть
|
||||
// Шаг 4: сравниваем current с desired
|
||||
currentSet := map[string]struct{}{}
|
||||
for _, p := range strings.Split(currentVal, ",") {
|
||||
if t := strings.TrimSpace(p); t != "" {
|
||||
currentSet[t] = struct{}{}
|
||||
}
|
||||
}
|
||||
newVal := currentVal + "," + ns
|
||||
|
||||
// Патчим все Fission deployments одним и тем же значением.
|
||||
// StrategicMergePatch обновляет только указанные поля (env var), не затрагивая остальные.
|
||||
// Обычный MergePatch заменил бы весь массив containers — нельзя использовать.
|
||||
// Проверяем симметричную разницу
|
||||
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": containerName,
|
||||
"name": "",
|
||||
"env": []any{
|
||||
map[string]any{
|
||||
"name": "FISSION_RESOURCE_NAMESPACES",
|
||||
@@ -1278,23 +1396,76 @@ func (s *server) addNSToFission(ctx context.Context, ns string) error {
|
||||
},
|
||||
},
|
||||
}
|
||||
patchBytes, err := json.Marshal(patch)
|
||||
if err != nil {
|
||||
return fmt.Errorf("marshal patch: %w", err)
|
||||
}
|
||||
for _, dep := range fissionDeployments {
|
||||
// В Fission каждый deployment имеет один контейнер с тем же именем что и deployment.
|
||||
// Подставляем имя контейнера под конкретный deployment для корректного merge key.
|
||||
patch["spec"].(map[string]any)["template"].(map[string]any)["spec"].(map[string]any)["containers"].([]any)[0].(map[string]any)["name"] = dep
|
||||
patchBytes, _ = json.Marshal(patch)
|
||||
_, patchErr := s.dyn.Resource(deploymentGVR).Namespace(fissionNS).Patch(
|
||||
ctx, dep, types.StrategicMergePatchType, patchBytes, metav1.PatchOptions{})
|
||||
if patchErr != nil {
|
||||
log.Printf("addNSToFission: patch deployment %s: %v", dep, patchErr)
|
||||
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("addNSToFission: added %s, new list: %s", ns, newVal)
|
||||
return nil
|
||||
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 {
|
||||
@@ -1354,19 +1525,40 @@ func (s *server) ensureUserNamespace(ctx context.Context, ns string) error {
|
||||
// даёт полный доступ ТОЛЬКО внутри конкретного namespace — это безопасно.
|
||||
//
|
||||
// Операция idempotent: если RoleBinding уже существует — IsAlreadyExists игнорируется.
|
||||
fissionSAs := []string{"fission-executor", "fission-router", "fission-buildermgr", "fission-kubewatcher", "fission-timer", "fission-fetcher", "fission-builder"}
|
||||
// Важно: 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": "fission-" + sa + "-user-ns",
|
||||
"name": sa.binding,
|
||||
"namespace": ns,
|
||||
},
|
||||
"roleRef": map[string]any{
|
||||
@@ -1377,25 +1569,150 @@ func (s *server) ensureUserNamespace(ctx context.Context, ns string) error {
|
||||
"subjects": []any{
|
||||
map[string]any{
|
||||
"kind": "ServiceAccount",
|
||||
"name": sa,
|
||||
"namespace": fissionSysNS,
|
||||
"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: %v", ns, sa, rbErr)
|
||||
log.Printf("ensureUserNamespace: create rolebinding %s/%s@%s: %v", ns, sa.name, subjectNS, rbErr)
|
||||
}
|
||||
}
|
||||
|
||||
// 1c. Регистрируем новый namespace в Fission (FISSION_RESOURCE_NAMESPACES).
|
||||
// Только при первом создании — повторный патч не нужен, Fission уже знает о namespace.
|
||||
// addNSToFission читает текущее значение переменной у router-а, добавляет ns и патчит
|
||||
// все Fission deployments (router, executor, buildermgr, kubewatcher, timer).
|
||||
// 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 {
|
||||
if patchErr := s.addNSToFission(ctx, ns); patchErr != nil {
|
||||
log.Printf("ensureUserNamespace: addNSToFission: %v", patchErr)
|
||||
select {
|
||||
case s.nsReconcileCh <- struct{}{}:
|
||||
default: // уже есть сигнал в буфере — не блокируем
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1472,10 +1789,10 @@ func (s *server) handleAuth(w http.ResponseWriter, r *http.Request) {
|
||||
ns = s.ns
|
||||
}
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second)
|
||||
ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second)
|
||||
defer cancel()
|
||||
if ensureErr := s.ensureUserNamespace(ctx, ns); ensureErr != nil {
|
||||
log.Printf("handleAuth: ensureUserNamespace %s: %v", ns, ensureErr)
|
||||
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})
|
||||
@@ -1520,6 +1837,7 @@ func (s *server) handleDeleteFunction(w http.ResponseWriter, r *http.Request, na
|
||||
}
|
||||
}
|
||||
|
||||
// Если environment больше не используется ни одной функцией — удаляем его.
|
||||
// Если environment больше не используется ни одной функцией — удаляем его.
|
||||
// Fission увидит удаление Environment CRD и убьёт pool deployment → поды умирают.
|
||||
if envName != "" {
|
||||
@@ -1528,6 +1846,13 @@ func (s *server) handleDeleteFunction(w http.ResponseWriter, r *http.Request, na
|
||||
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})
|
||||
}
|
||||
|
||||
|
||||
+22
-3
@@ -457,7 +457,8 @@
|
||||
<button class="btn ghost" onclick="closeInvoke()">Отмена</button>
|
||||
<button id="i-submit" class="btn" onclick="submitInvoke()">Вызвать</button>
|
||||
</div>
|
||||
<div style="margin-top:10px;">
|
||||
<div id="i-status" style="margin-top:8px;font-size:13px;color:var(--text-secondary);min-height:20px;"></div>
|
||||
<div style="margin-top:6px;">
|
||||
<label>Response</label>
|
||||
<textarea id="i-resp" readonly style="min-height:160px;"></textarea>
|
||||
</div>
|
||||
@@ -696,16 +697,34 @@
|
||||
async function submitInvoke() {
|
||||
if (!S.currentInvoke) return;
|
||||
const btn = document.getElementById('i-submit');
|
||||
const statusEl = document.getElementById('i-status');
|
||||
const respEl = document.getElementById('i-resp');
|
||||
btn.disabled = true;
|
||||
respEl.value = '';
|
||||
let elapsed = 0;
|
||||
statusEl.textContent = 'Вызов...';
|
||||
const timer = setInterval(() => {
|
||||
elapsed++;
|
||||
if (elapsed < 5) {
|
||||
statusEl.textContent = 'Вызов... ' + elapsed + 'с';
|
||||
} else if (elapsed < 10) {
|
||||
statusEl.textContent = '⏳ Холодный старт — прогрев пула... ' + elapsed + 'с';
|
||||
} else {
|
||||
statusEl.textContent = '⏳ Холодный старт — ещё немного... ' + elapsed + 'с';
|
||||
}
|
||||
}, 1000);
|
||||
try {
|
||||
const raw = document.getElementById('i-body').value.trim();
|
||||
let parsed = {};
|
||||
if (raw) parsed = JSON.parse(raw);
|
||||
const result = await requestJSON(API_BASE + '/functions/' + encodeURIComponent(S.currentInvoke) + '/invoke', 'POST', parsed);
|
||||
document.getElementById('i-resp').value = JSON.stringify(result, null, 2);
|
||||
statusEl.textContent = '✓ Выполнено за ' + elapsed + 'с';
|
||||
respEl.value = JSON.stringify(result, null, 2);
|
||||
} catch (e) {
|
||||
document.getElementById('i-resp').value = '\u041e\u0448\u0438\u0431\u043a\u0430 \u0432\u044b\u0437\u043e\u0432\u0430: ' + e.message;
|
||||
statusEl.textContent = '✗ Ошибка после ' + elapsed + 'с';
|
||||
respEl.value = 'Ошибка вызова: ' + e.message;
|
||||
} finally {
|
||||
clearInterval(timer);
|
||||
btn.disabled = false;
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user