902 lines
27 KiB
Go
902 lines
27 KiB
Go
package main
|
|
|
|
import (
|
|
"archive/zip"
|
|
"bytes"
|
|
"context"
|
|
"encoding/base64"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"log"
|
|
"net"
|
|
"net/http"
|
|
"os"
|
|
"sort"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
"unicode/utf8"
|
|
|
|
"fission-console/ui"
|
|
|
|
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
|
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
|
|
"k8s.io/apimachinery/pkg/runtime/schema"
|
|
"k8s.io/client-go/dynamic"
|
|
"k8s.io/client-go/rest"
|
|
"k8s.io/client-go/tools/clientcmd"
|
|
)
|
|
|
|
var (
|
|
environmentGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "environments"}
|
|
packageGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "packages"}
|
|
functionGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "functions"}
|
|
httpTrigGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "httptriggers"}
|
|
timeTrigGVR = schema.GroupVersionResource{Group: "fission.io", Version: "v1", Resource: "timetriggers"}
|
|
)
|
|
|
|
const defaultSATokenPath = "/var/run/secrets/kubernetes.io/serviceaccount/token"
|
|
|
|
type server struct {
|
|
dyn dynamic.Interface
|
|
ns string
|
|
routerURL string
|
|
http *http.Client
|
|
saTokenPath string
|
|
invokeTimeout time.Duration
|
|
|
|
authUser string
|
|
authPass string
|
|
|
|
tokenMu sync.Mutex
|
|
cachedJWT string
|
|
tokenExpAt time.Time
|
|
}
|
|
|
|
type createFunctionRequest struct {
|
|
Name string `json:"name"`
|
|
Environment string `json:"environment"`
|
|
Code string `json:"code"`
|
|
Entrypoint string `json:"entrypoint"`
|
|
Route string `json:"route"`
|
|
Methods []string `json:"methods"`
|
|
}
|
|
|
|
type updateCodeRequest struct {
|
|
Code string `json:"code"`
|
|
}
|
|
|
|
func main() {
|
|
kubeconfig := strings.TrimSpace(os.Getenv("KUBECONFIG"))
|
|
namespace := envDefault("FISSION_NAMESPACE", "default")
|
|
routerURL := strings.TrimRight(envDefault("FISSION_ROUTER_URL", "http://router.fission.svc.cluster.local"), "/")
|
|
port := envDefault("PORT", "8090")
|
|
httpTimeout := envDurationDefault("FISSION_HTTP_TIMEOUT", 30*time.Second)
|
|
invokeTimeout := envDurationDefault("FISSION_INVOKE_TIMEOUT", 20*time.Second)
|
|
|
|
cfg, err := buildConfig(kubeconfig)
|
|
if err != nil {
|
|
log.Fatalf("build kube config: %v", err)
|
|
}
|
|
|
|
dyn, err := dynamic.NewForConfig(cfg)
|
|
if err != nil {
|
|
log.Fatalf("create dynamic client: %v", err)
|
|
}
|
|
|
|
authUser := envDefault("FISSION_AUTH_USERNAME", "")
|
|
authPass := envDefault("FISSION_AUTH_PASSWORD", "")
|
|
saTokenPath := envDefault("SA_TOKEN_PATH", defaultSATokenPath)
|
|
|
|
s := &server{
|
|
dyn: dyn,
|
|
ns: namespace,
|
|
routerURL: routerURL,
|
|
http: &http.Client{Timeout: httpTimeout},
|
|
saTokenPath: saTokenPath,
|
|
invokeTimeout: invokeTimeout,
|
|
authUser: authUser,
|
|
authPass: authPass,
|
|
}
|
|
|
|
mux := http.NewServeMux()
|
|
mux.HandleFunc("/health", func(w http.ResponseWriter, _ *http.Request) {
|
|
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
|
|
_, _ = w.Write([]byte("ok\n"))
|
|
})
|
|
mux.HandleFunc("/console/health", func(w http.ResponseWriter, _ *http.Request) {
|
|
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
|
|
_, _ = w.Write([]byte("ok\n"))
|
|
})
|
|
|
|
uiHandler := http.StripPrefix("/console", ui.Handler())
|
|
mux.Handle("/console", uiHandler)
|
|
mux.Handle("/console/", uiHandler)
|
|
|
|
mux.HandleFunc("/api/environments", s.handleList(environmentGVR))
|
|
mux.HandleFunc("/api/packages", s.handleList(packageGVR))
|
|
mux.HandleFunc("/api/functions", s.handleFunctionsRoot)
|
|
mux.HandleFunc("/api/functions/", s.handleFunctionsAction)
|
|
mux.HandleFunc("/api/httptriggers", s.handleList(httpTrigGVR))
|
|
mux.HandleFunc("/api/timetriggers", s.handleList(timeTrigGVR))
|
|
|
|
mux.HandleFunc("/console/api/environments", s.handleList(environmentGVR))
|
|
mux.HandleFunc("/console/api/packages", s.handleList(packageGVR))
|
|
mux.HandleFunc("/console/api/functions", s.handleFunctionsRoot)
|
|
mux.HandleFunc("/console/api/functions/", s.handleFunctionsAction)
|
|
mux.HandleFunc("/console/api/httptriggers", s.handleList(httpTrigGVR))
|
|
mux.HandleFunc("/console/api/timetriggers", s.handleList(timeTrigGVR))
|
|
|
|
httpServer := &http.Server{
|
|
Addr: ":" + port,
|
|
Handler: withSecurityHeaders(withCORS(logRequests(mux))),
|
|
ReadHeaderTimeout: 10 * time.Second,
|
|
}
|
|
|
|
log.Printf("fission-console listening on :%s (namespace=%s)", port, namespace)
|
|
log.Fatal(httpServer.ListenAndServe())
|
|
}
|
|
|
|
func (s *server) handleFunctionsRoot(w http.ResponseWriter, r *http.Request) {
|
|
switch r.Method {
|
|
case http.MethodGet:
|
|
s.handleList(functionGVR)(w, r)
|
|
case http.MethodPost:
|
|
s.handleCreateFunction(w, r)
|
|
default:
|
|
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
|
}
|
|
}
|
|
|
|
func (s *server) handleFunctionsAction(w http.ResponseWriter, r *http.Request) {
|
|
path := strings.TrimPrefix(r.URL.Path, "/api/functions/")
|
|
if path == r.URL.Path {
|
|
path = strings.TrimPrefix(r.URL.Path, "/console/api/functions/")
|
|
}
|
|
path = strings.Trim(path, "/")
|
|
if path == "" {
|
|
http.NotFound(w, r)
|
|
return
|
|
}
|
|
|
|
parts := strings.Split(path, "/")
|
|
name := strings.TrimSpace(parts[0])
|
|
if name == "" {
|
|
writeJSONError(w, http.StatusBadRequest, "function name is required")
|
|
return
|
|
}
|
|
|
|
if len(parts) == 1 {
|
|
switch r.Method {
|
|
case http.MethodGet:
|
|
s.handleGetFunction(w, r, name)
|
|
case http.MethodDelete:
|
|
s.handleDeleteFunction(w, r, name)
|
|
default:
|
|
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
|
}
|
|
return
|
|
}
|
|
|
|
if len(parts) == 2 && parts[1] == "code" && r.Method == http.MethodPut {
|
|
s.handleUpdateFunctionCode(w, r, name)
|
|
return
|
|
}
|
|
|
|
if len(parts) == 2 && parts[1] == "invoke" && r.Method == http.MethodPost {
|
|
s.handleInvokeFunction(w, r, name)
|
|
return
|
|
}
|
|
|
|
http.NotFound(w, r)
|
|
}
|
|
|
|
func (s *server) handleCreateFunction(w http.ResponseWriter, r *http.Request) {
|
|
var req createFunctionRequest
|
|
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
|
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("decode request: %v", err))
|
|
return
|
|
}
|
|
|
|
req.Name = strings.TrimSpace(req.Name)
|
|
req.Environment = strings.TrimSpace(req.Environment)
|
|
req.Code = strings.TrimSpace(req.Code)
|
|
req.Entrypoint = strings.TrimSpace(req.Entrypoint)
|
|
req.Route = strings.TrimSpace(req.Route)
|
|
|
|
if req.Name == "" || req.Environment == "" || req.Code == "" {
|
|
writeJSONError(w, http.StatusBadRequest, "name, environment and code are required")
|
|
return
|
|
}
|
|
if req.Entrypoint == "" {
|
|
req.Entrypoint = "main.main"
|
|
}
|
|
if req.Route == "" {
|
|
req.Route = "/" + req.Name
|
|
}
|
|
if !strings.HasPrefix(req.Route, "/") {
|
|
req.Route = "/" + req.Route
|
|
}
|
|
req.Methods = normalizeMethods(req.Methods)
|
|
|
|
ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second)
|
|
defer cancel()
|
|
|
|
if _, err := s.dyn.Resource(environmentGVR).Namespace(s.ns).Get(ctx, req.Environment, metav1.GetOptions{}); err != nil {
|
|
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("environment %q not found: %v", req.Environment, err))
|
|
return
|
|
}
|
|
|
|
pkgName := req.Name + "-pkg"
|
|
triggerName := req.Name + "-route"
|
|
methodValues := make([]any, 0, len(req.Methods))
|
|
for _, method := range req.Methods {
|
|
methodValues = append(methodValues, method)
|
|
}
|
|
|
|
literal := base64.StdEncoding.EncodeToString([]byte(req.Code))
|
|
pkg := &unstructured.Unstructured{Object: map[string]any{
|
|
"apiVersion": "fission.io/v1",
|
|
"kind": "Package",
|
|
"metadata": map[string]any{
|
|
"name": pkgName,
|
|
"namespace": s.ns,
|
|
},
|
|
"spec": map[string]any{
|
|
"deployment": map[string]any{
|
|
"type": "literal",
|
|
"literal": literal,
|
|
},
|
|
"environment": map[string]any{
|
|
"name": req.Environment,
|
|
"namespace": s.ns,
|
|
},
|
|
"source": map[string]any{},
|
|
},
|
|
}}
|
|
|
|
if _, err := s.dyn.Resource(packageGVR).Namespace(s.ns).Create(ctx, pkg, metav1.CreateOptions{}); err != nil {
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create package: %v", err))
|
|
return
|
|
}
|
|
|
|
fn := &unstructured.Unstructured{Object: map[string]any{
|
|
"apiVersion": "fission.io/v1",
|
|
"kind": "Function",
|
|
"metadata": map[string]any{
|
|
"name": req.Name,
|
|
"namespace": s.ns,
|
|
},
|
|
"spec": map[string]any{
|
|
"environment": map[string]any{
|
|
"name": req.Environment,
|
|
"namespace": s.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},
|
|
"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{})
|
|
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": s.ns,
|
|
},
|
|
"spec": map[string]any{
|
|
"relativeurl": req.Route,
|
|
"methods": methodValues,
|
|
"createingress": true,
|
|
"functionref": map[string]any{
|
|
"type": "name",
|
|
"name": req.Name,
|
|
},
|
|
},
|
|
}}
|
|
|
|
if _, err := s.dyn.Resource(httpTrigGVR).Namespace(s.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{})
|
|
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,
|
|
})
|
|
}
|
|
|
|
func (s *server) handleGetFunction(w http.ResponseWriter, r *http.Request, name string) {
|
|
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second)
|
|
defer cancel()
|
|
|
|
fn, err := s.dyn.Resource(functionGVR).Namespace(s.ns).Get(ctx, name, metav1.GetOptions{})
|
|
if err != nil {
|
|
status := http.StatusBadGateway
|
|
if apierrors.IsNotFound(err) {
|
|
status = http.StatusNotFound
|
|
}
|
|
writeJSONError(w, status, fmt.Sprintf("get function %q: %v", name, err))
|
|
return
|
|
}
|
|
|
|
packageName, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name")
|
|
environment, _, _ := unstructured.NestedString(fn.Object, "spec", "environment", "name")
|
|
entrypoint, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "functionName")
|
|
|
|
code := ""
|
|
if packageName != "" {
|
|
pkg, pkgErr := s.dyn.Resource(packageGVR).Namespace(s.ns).Get(ctx, packageName, metav1.GetOptions{})
|
|
if pkgErr == nil {
|
|
code = s.extractPackageSourceCode(ctx, pkg)
|
|
}
|
|
}
|
|
|
|
route := ""
|
|
methods := []string{}
|
|
triggers, trigErr := s.dyn.Resource(httpTrigGVR).Namespace(s.ns).List(ctx, metav1.ListOptions{})
|
|
if trigErr == nil {
|
|
for _, trig := range triggers.Items {
|
|
refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name")
|
|
if refName != name {
|
|
continue
|
|
}
|
|
route, _, _ = unstructured.NestedString(trig.Object, "spec", "relativeurl")
|
|
methods, _, _ = unstructured.NestedStringSlice(trig.Object, "spec", "methods")
|
|
break
|
|
}
|
|
}
|
|
|
|
writeAnyJSON(w, http.StatusOK, map[string]any{
|
|
"name": name,
|
|
"namespace": s.ns,
|
|
"environment": environment,
|
|
"package": packageName,
|
|
"entrypoint": entrypoint,
|
|
"code": code,
|
|
"route": route,
|
|
"methods": methods,
|
|
"raw": fn.Object,
|
|
})
|
|
}
|
|
|
|
func (s *server) handleUpdateFunctionCode(w http.ResponseWriter, r *http.Request, name string) {
|
|
var req updateCodeRequest
|
|
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
|
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("decode request: %v", err))
|
|
return
|
|
}
|
|
req.Code = strings.TrimSpace(req.Code)
|
|
if req.Code == "" {
|
|
writeJSONError(w, http.StatusBadRequest, "code is required")
|
|
return
|
|
}
|
|
|
|
ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second)
|
|
defer cancel()
|
|
|
|
fn, err := s.dyn.Resource(functionGVR).Namespace(s.ns).Get(ctx, name, metav1.GetOptions{})
|
|
if err != nil {
|
|
status := http.StatusBadGateway
|
|
if apierrors.IsNotFound(err) {
|
|
status = http.StatusNotFound
|
|
}
|
|
writeJSONError(w, status, fmt.Sprintf("get function %q: %v", name, err))
|
|
return
|
|
}
|
|
|
|
pkgName, _, _ := unstructured.NestedString(fn.Object, "spec", "package", "packageref", "name")
|
|
if pkgName == "" {
|
|
writeJSONError(w, http.StatusBadGateway, "function has no package reference")
|
|
return
|
|
}
|
|
|
|
pkg, err := s.dyn.Resource(packageGVR).Namespace(s.ns).Get(ctx, pkgName, metav1.GetOptions{})
|
|
if err != nil {
|
|
status := http.StatusBadGateway
|
|
if apierrors.IsNotFound(err) {
|
|
status = http.StatusNotFound
|
|
}
|
|
writeJSONError(w, status, fmt.Sprintf("get package %q: %v", pkgName, err))
|
|
return
|
|
}
|
|
|
|
literal := base64.StdEncoding.EncodeToString([]byte(req.Code))
|
|
if err := unstructured.SetNestedField(pkg.Object, literal, "spec", "deployment", "literal"); err != nil {
|
|
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set package literal: %v", err))
|
|
return
|
|
}
|
|
|
|
if _, err := s.dyn.Resource(packageGVR).Namespace(s.ns).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{})
|
|
if err != nil {
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get updated package %q: %v", pkgName, err))
|
|
return
|
|
}
|
|
|
|
if err := unstructured.SetNestedField(fn.Object, updatedPkg.GetResourceVersion(), "spec", "package", "packageref", "resourceversion"); err != nil {
|
|
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("set function package resourceversion: %v", err))
|
|
return
|
|
}
|
|
|
|
if _, err := s.dyn.Resource(functionGVR).Namespace(s.ns).Update(ctx, fn, metav1.UpdateOptions{}); err != nil {
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("update function %q package ref: %v", name, err))
|
|
return
|
|
}
|
|
|
|
writeAnyJSON(w, http.StatusOK, map[string]any{"updated": true, "package": pkgName, "package_resourceversion": updatedPkg.GetResourceVersion()})
|
|
}
|
|
|
|
func (s *server) handleInvokeFunction(w http.ResponseWriter, r *http.Request, name string) {
|
|
bodyBytes, err := io.ReadAll(r.Body)
|
|
if err != nil {
|
|
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("read request body: %v", err))
|
|
return
|
|
}
|
|
if len(bytes.TrimSpace(bodyBytes)) == 0 {
|
|
bodyBytes = []byte("{}")
|
|
}
|
|
|
|
invokeTimeout := s.invokeTimeout
|
|
if invokeTimeout <= 0 {
|
|
invokeTimeout = 20 * time.Second
|
|
}
|
|
|
|
ctx, cancel := context.WithTimeout(r.Context(), invokeTimeout)
|
|
defer cancel()
|
|
|
|
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{})
|
|
if err == nil {
|
|
for _, trig := range triggers.Items {
|
|
refName, _, _ := unstructured.NestedString(trig.Object, "spec", "functionref", "name")
|
|
if refName != name {
|
|
continue
|
|
}
|
|
route, _, _ := unstructured.NestedString(trig.Object, "spec", "relativeurl")
|
|
methods, _, _ := unstructured.NestedStringSlice(trig.Object, "spec", "methods")
|
|
hasPost := false
|
|
hasGet := false
|
|
for _, method := range methods {
|
|
m := strings.ToUpper(strings.TrimSpace(method))
|
|
if m == http.MethodPost {
|
|
hasPost = true
|
|
}
|
|
if m == http.MethodGet {
|
|
hasGet = true
|
|
}
|
|
}
|
|
if route != "" {
|
|
if !strings.HasPrefix(route, "/") {
|
|
route = "/" + route
|
|
}
|
|
invokeURL = s.routerURL + route
|
|
if !hasPost && hasGet {
|
|
invokeMethod = http.MethodGet
|
|
}
|
|
break
|
|
}
|
|
}
|
|
}
|
|
|
|
start := time.Now()
|
|
var invokeBody io.Reader
|
|
if invokeMethod == http.MethodPost {
|
|
invokeBody = bytes.NewReader(bodyBytes)
|
|
}
|
|
req, err := http.NewRequestWithContext(ctx, invokeMethod, invokeURL, invokeBody)
|
|
if err != nil {
|
|
writeJSONError(w, http.StatusInternalServerError, fmt.Sprintf("build invoke request: %v", err))
|
|
return
|
|
}
|
|
if invokeMethod == http.MethodPost {
|
|
req.Header.Set("Content-Type", "application/json")
|
|
}
|
|
if token := s.getRouterToken(); token != "" {
|
|
req.Header.Set("Authorization", "Bearer "+token)
|
|
}
|
|
|
|
resp, err := s.http.Do(req)
|
|
if err != nil {
|
|
if errors.Is(err, context.DeadlineExceeded) {
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q timeout after %s: function specialization likely failed (for example, syntax error)", name, invokeTimeout))
|
|
return
|
|
}
|
|
var netErr net.Error
|
|
if errors.As(err, &netErr) && netErr.Timeout() {
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q timeout after %s: function specialization likely failed (for example, syntax error)", name, invokeTimeout))
|
|
return
|
|
}
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("invoke %q: %v", name, err))
|
|
return
|
|
}
|
|
defer resp.Body.Close()
|
|
|
|
respBody, _ := io.ReadAll(resp.Body)
|
|
writeAnyJSON(w, http.StatusOK, map[string]any{
|
|
"status": resp.StatusCode,
|
|
"latency_ms": time.Since(start).Milliseconds(),
|
|
"invoke_url": invokeURL,
|
|
"response_raw": string(respBody),
|
|
})
|
|
}
|
|
|
|
func (s *server) readSAToken() string {
|
|
if s.saTokenPath == "" {
|
|
return ""
|
|
}
|
|
data, err := os.ReadFile(s.saTokenPath)
|
|
if err != nil {
|
|
return ""
|
|
}
|
|
return strings.TrimSpace(string(data))
|
|
}
|
|
|
|
func (s *server) getRouterToken() string {
|
|
if s.authUser == "" || s.authPass == "" {
|
|
return s.readSAToken()
|
|
}
|
|
|
|
s.tokenMu.Lock()
|
|
defer s.tokenMu.Unlock()
|
|
|
|
if s.cachedJWT != "" && time.Now().Before(s.tokenExpAt) {
|
|
return s.cachedJWT
|
|
}
|
|
|
|
loginURL := s.routerURL + "/auth/login"
|
|
body, _ := json.Marshal(map[string]string{"username": s.authUser, "password": s.authPass})
|
|
resp, err := s.http.Post(loginURL, "application/json", bytes.NewReader(body))
|
|
if err != nil {
|
|
log.Printf("router login failed: %v", err)
|
|
return s.readSAToken()
|
|
}
|
|
defer resp.Body.Close()
|
|
|
|
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated {
|
|
respBody, _ := io.ReadAll(resp.Body)
|
|
log.Printf("router login %d: %s", resp.StatusCode, string(respBody))
|
|
return s.readSAToken()
|
|
}
|
|
|
|
var result struct {
|
|
AccessToken string `json:"accesstoken"`
|
|
}
|
|
if err := json.NewDecoder(resp.Body).Decode(&result); err != nil || result.AccessToken == "" {
|
|
log.Printf("router login decode error: %v", err)
|
|
return s.readSAToken()
|
|
}
|
|
|
|
s.cachedJWT = result.AccessToken
|
|
s.tokenExpAt = time.Now().Add(100 * time.Second)
|
|
log.Printf("router JWT obtained, expires in 100s")
|
|
return s.cachedJWT
|
|
}
|
|
|
|
func (s *server) handleDeleteFunction(w http.ResponseWriter, r *http.Request, name string) {
|
|
ctx, cancel := context.WithTimeout(r.Context(), 20*time.Second)
|
|
defer cancel()
|
|
|
|
var pkgName string
|
|
fn, err := s.dyn.Resource(functionGVR).Namespace(s.ns).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{})
|
|
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{})
|
|
}
|
|
}
|
|
}
|
|
|
|
if err := s.dyn.Resource(functionGVR).Namespace(s.ns).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) {
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("delete package %q: %v", pkgName, err))
|
|
return
|
|
}
|
|
}
|
|
|
|
writeAnyJSON(w, http.StatusOK, map[string]any{"deleted": true, "name": name, "package": pkgName})
|
|
}
|
|
|
|
func buildConfig(kubeconfig string) (*rest.Config, error) {
|
|
if kubeconfig != "" {
|
|
cfg, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
|
|
if err == nil {
|
|
return cfg, nil
|
|
}
|
|
return nil, fmt.Errorf("kubeconfig %s: %w", kubeconfig, err)
|
|
}
|
|
|
|
cfg, err := rest.InClusterConfig()
|
|
if err == nil {
|
|
return cfg, nil
|
|
}
|
|
|
|
loadingRules := &clientcmd.ClientConfigLoadingRules{}
|
|
clientCfg := clientcmd.NewNonInteractiveDeferredLoadingClientConfig(loadingRules, &clientcmd.ConfigOverrides{})
|
|
return clientCfg.ClientConfig()
|
|
}
|
|
|
|
func (s *server) handleList(gvr schema.GroupVersionResource) http.HandlerFunc {
|
|
return func(w http.ResponseWriter, r *http.Request) {
|
|
if r.Method != http.MethodGet {
|
|
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
|
return
|
|
}
|
|
|
|
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second)
|
|
defer cancel()
|
|
|
|
list, err := s.dyn.Resource(gvr).Namespace(s.ns).List(ctx, metav1.ListOptions{})
|
|
if err != nil {
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("list %s: %v", gvr.Resource, err))
|
|
return
|
|
}
|
|
|
|
writeJSON(w, http.StatusOK, list.Items)
|
|
}
|
|
}
|
|
|
|
func writeJSON(w http.ResponseWriter, status int, data []unstructured.Unstructured) {
|
|
w.Header().Set("Content-Type", "application/json; charset=utf-8")
|
|
w.WriteHeader(status)
|
|
_ = json.NewEncoder(w).Encode(data)
|
|
}
|
|
|
|
func writeAnyJSON(w http.ResponseWriter, status int, data any) {
|
|
w.Header().Set("Content-Type", "application/json; charset=utf-8")
|
|
w.WriteHeader(status)
|
|
_ = json.NewEncoder(w).Encode(data)
|
|
}
|
|
|
|
func writeJSONError(w http.ResponseWriter, status int, msg string) {
|
|
w.Header().Set("Content-Type", "application/json; charset=utf-8")
|
|
w.WriteHeader(status)
|
|
_ = json.NewEncoder(w).Encode(map[string]any{"error": msg})
|
|
}
|
|
|
|
func logRequests(next http.Handler) http.Handler {
|
|
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
log.Printf("%s %s", r.Method, r.URL.Path)
|
|
next.ServeHTTP(w, r)
|
|
})
|
|
}
|
|
|
|
func withCORS(next http.Handler) http.Handler {
|
|
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
w.Header().Set("Access-Control-Allow-Origin", "*")
|
|
w.Header().Set("Access-Control-Allow-Methods", "GET,POST,PUT,PATCH,DELETE,OPTIONS")
|
|
w.Header().Set("Access-Control-Allow-Headers", "Content-Type, Authorization")
|
|
if r.Method == http.MethodOptions {
|
|
w.WriteHeader(http.StatusNoContent)
|
|
return
|
|
}
|
|
next.ServeHTTP(w, r)
|
|
})
|
|
}
|
|
|
|
func withSecurityHeaders(next http.Handler) http.Handler {
|
|
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
w.Header().Set("X-Content-Type-Options", "nosniff")
|
|
w.Header().Set("X-Frame-Options", "DENY")
|
|
w.Header().Set("Referrer-Policy", "strict-origin-when-cross-origin")
|
|
w.Header().Set("Permissions-Policy", "camera=(), microphone=(), geolocation=()")
|
|
w.Header().Set("Content-Security-Policy", "default-src 'self'; script-src 'self' 'unsafe-inline'; style-src 'self' 'unsafe-inline'; img-src 'self' data:; connect-src 'self'; font-src 'self' data:; object-src 'none'; frame-ancestors 'none'; base-uri 'self'; form-action 'self'; upgrade-insecure-requests; block-all-mixed-content")
|
|
next.ServeHTTP(w, r)
|
|
})
|
|
}
|
|
|
|
func envDefault(key, fallback string) string {
|
|
if v := strings.TrimSpace(os.Getenv(key)); v != "" {
|
|
return v
|
|
}
|
|
return fallback
|
|
}
|
|
|
|
func envDurationDefault(key string, fallback time.Duration) time.Duration {
|
|
raw := strings.TrimSpace(os.Getenv(key))
|
|
if raw == "" {
|
|
return fallback
|
|
}
|
|
d, err := time.ParseDuration(raw)
|
|
if err != nil || d <= 0 {
|
|
log.Printf("invalid duration for %s=%q, using default %s", key, raw, fallback)
|
|
return fallback
|
|
}
|
|
return d
|
|
}
|
|
|
|
func normalizeMethods(in []string) []string {
|
|
if len(in) == 0 {
|
|
return []string{"GET"}
|
|
}
|
|
out := make([]string, 0, len(in))
|
|
seen := map[string]bool{}
|
|
for _, method := range in {
|
|
m := strings.ToUpper(strings.TrimSpace(method))
|
|
if m == "" || seen[m] {
|
|
continue
|
|
}
|
|
seen[m] = true
|
|
out = append(out, m)
|
|
}
|
|
if len(out) == 0 {
|
|
return []string{"GET"}
|
|
}
|
|
return out
|
|
}
|
|
|
|
func (s *server) extractPackageSourceCode(ctx context.Context, pkg *unstructured.Unstructured) string {
|
|
literalPaths := [][]string{
|
|
{"spec", "source", "literal"},
|
|
{"spec", "deployment", "literal"},
|
|
}
|
|
for _, p := range literalPaths {
|
|
literal, found, _ := unstructured.NestedString(pkg.Object, p...)
|
|
if !found || strings.TrimSpace(literal) == "" {
|
|
continue
|
|
}
|
|
if decodedCode, decErr := decodeLiteralToSource(literal); decErr == nil && strings.TrimSpace(decodedCode) != "" {
|
|
return decodedCode
|
|
}
|
|
}
|
|
|
|
urlPaths := [][]string{
|
|
{"spec", "source", "url"},
|
|
{"spec", "deployment", "url"},
|
|
}
|
|
for _, p := range urlPaths {
|
|
urlValue, found, _ := unstructured.NestedString(pkg.Object, p...)
|
|
if !found || strings.TrimSpace(urlValue) == "" {
|
|
continue
|
|
}
|
|
|
|
archiveBytes, fetchErr := s.fetchPackageArchive(ctx, urlValue)
|
|
if fetchErr != nil {
|
|
continue
|
|
}
|
|
|
|
decodedCode, decErr := decodeArchiveBytesToSource(archiveBytes)
|
|
if decErr == nil && strings.TrimSpace(decodedCode) != "" {
|
|
return decodedCode
|
|
}
|
|
}
|
|
|
|
return ""
|
|
}
|
|
|
|
func (s *server) fetchPackageArchive(ctx context.Context, archiveURL string) ([]byte, error) {
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, archiveURL, nil)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
resp, err := s.http.Do(req)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer resp.Body.Close()
|
|
|
|
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
|
return nil, fmt.Errorf("archive request failed: %s", resp.Status)
|
|
}
|
|
|
|
return io.ReadAll(resp.Body)
|
|
}
|
|
|
|
func decodeLiteralToSource(literal string) (string, error) {
|
|
decoded, err := base64.StdEncoding.DecodeString(literal)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
|
|
return decodeArchiveBytesToSource(decoded)
|
|
}
|
|
|
|
func decodeArchiveBytesToSource(decoded []byte) (string, error) {
|
|
if len(decoded) == 0 {
|
|
return "", fmt.Errorf("empty payload")
|
|
}
|
|
|
|
if utf8.Valid(decoded) {
|
|
return string(decoded), nil
|
|
}
|
|
|
|
if len(decoded) >= 4 && bytes.Equal(decoded[:4], []byte{'P', 'K', 3, 4}) {
|
|
if src, zipErr := decodeZipSource(decoded); zipErr == nil {
|
|
return src, nil
|
|
}
|
|
}
|
|
|
|
return "", fmt.Errorf("payload does not contain utf-8 source")
|
|
}
|
|
|
|
func decodeZipSource(zipBytes []byte) (string, error) {
|
|
reader, err := zip.NewReader(bytes.NewReader(zipBytes), int64(len(zipBytes)))
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
|
|
preferred := []string{"main.py", "main.js", "main.go"}
|
|
for _, name := range preferred {
|
|
for _, file := range reader.File {
|
|
if strings.EqualFold(file.Name, name) {
|
|
content, readErr := readZipFile(file)
|
|
if readErr != nil {
|
|
return "", readErr
|
|
}
|
|
if utf8.Valid(content) {
|
|
return string(content), nil
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
files := make([]*zip.File, 0, len(reader.File))
|
|
for _, file := range reader.File {
|
|
if file.FileInfo().IsDir() {
|
|
continue
|
|
}
|
|
files = append(files, file)
|
|
}
|
|
sort.Slice(files, func(i, j int) bool {
|
|
return files[i].Name < files[j].Name
|
|
})
|
|
|
|
for _, file := range files {
|
|
content, readErr := readZipFile(file)
|
|
if readErr != nil {
|
|
continue
|
|
}
|
|
if utf8.Valid(content) {
|
|
return string(content), nil
|
|
}
|
|
}
|
|
|
|
return "", fmt.Errorf("zip archive does not contain utf-8 source files")
|
|
}
|
|
|
|
func readZipFile(file *zip.File) ([]byte, error) {
|
|
rc, err := file.Open()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rc.Close()
|
|
|
|
return io.ReadAll(rc)
|
|
}
|