- internal/billing: Store interface + pgStore (pgx/v5 pool) + NoopStore - factory.go: NewStore() from BILLING_DSN env var (NoopStore fallback) - Server.billing: injected into Config, initialized in main - handleInvokeFunction: record invocation (TriggerConsole) - invokeInternalFunction: record invocation (TriggerHTTP) - handleCreateFunction: record event_type=create - handleCloneFunction: record event_type=clone - handleDeleteFunction: record event_type=delete - table: invocations (namespace, function_name, trigger_type, duration_ms, status_code...) - BILLING_DSN added to console.yaml - pgx/v5 added to go.mod/go.sum
334 lines
12 KiB
Go
334 lines
12 KiB
Go
// Package api — клонирование функций.
|
|
//
|
|
// handleCloneFunction создаёт полную копию функции с новым именем:
|
|
// - скачивает архив из storagesvc (или копирует literal)
|
|
// - заливает новый архив (отдельный объект в S3)
|
|
// - создаёт новый Package, Function и HTTPTrigger
|
|
//
|
|
// Архив переливается заново, чтобы удаление оригинала не сломало клон.
|
|
package api
|
|
|
|
import (
|
|
"context"
|
|
"encoding/base64"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"regexp"
|
|
"strings"
|
|
"time"
|
|
|
|
"fission-console/internal/billing"
|
|
"fission-console/internal/fission"
|
|
"fission-console/internal/runtime"
|
|
|
|
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
|
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
|
|
)
|
|
|
|
// handleCloneFunction — POST /functions/:name/clone
|
|
// Body: {"new_name": "my-clone", "route": "/my-clone"} (route необязателен)
|
|
func (s *Server) handleCloneFunction(w http.ResponseWriter, r *http.Request, srcName string) {
|
|
ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second)
|
|
defer cancel()
|
|
ns := s.userNS(r)
|
|
|
|
var req struct {
|
|
NewName string `json:"new_name"`
|
|
Route string `json:"route"`
|
|
}
|
|
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
|
writeJSONError(w, http.StatusBadRequest, fmt.Sprintf("decode request: %v", err))
|
|
return
|
|
}
|
|
req.NewName = strings.TrimSpace(req.NewName)
|
|
req.Route = strings.TrimSpace(req.Route)
|
|
|
|
// Валидация нового имени
|
|
if req.NewName == "" {
|
|
writeJSONError(w, http.StatusBadRequest, "new_name is required")
|
|
return
|
|
}
|
|
validName := regexp.MustCompile(`^[a-z0-9]([a-z0-9-]*[a-z0-9])?$`)
|
|
if !validName.MatchString(req.NewName) || len(req.NewName) > 57 {
|
|
writeJSONError(w, http.StatusBadRequest, "invalid new_name: must match ^[a-z0-9]([a-z0-9-]*[a-z0-9])?$ and be <= 57 chars")
|
|
return
|
|
}
|
|
|
|
// Получаем исходную функцию
|
|
srcFn, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Get(ctx, srcName, metav1.GetOptions{})
|
|
if err != nil {
|
|
if apierrors.IsNotFound(err) {
|
|
writeJSONError(w, http.StatusNotFound, fmt.Sprintf("function %q not found", srcName))
|
|
return
|
|
}
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get function: %v", err))
|
|
return
|
|
}
|
|
|
|
// Получаем имя пакета исходной функции
|
|
srcPkgName, _, _ := unstructured.NestedString(srcFn.Object, "spec", "package", "packageref", "name")
|
|
if srcPkgName == "" {
|
|
writeJSONError(w, http.StatusBadGateway, "source function has no package reference")
|
|
return
|
|
}
|
|
|
|
// Получаем исходный Package
|
|
srcPkg, err := s.dyn.Resource(fission.PackageGVR).Namespace(ns).Get(ctx, srcPkgName, metav1.GetOptions{})
|
|
if err != nil {
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("get source package: %v", err))
|
|
return
|
|
}
|
|
|
|
// Скачиваем байты архива из Package
|
|
archiveBytes, err := s.downloadPackageBytes(ctx, srcPkg)
|
|
if err != nil {
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("download archive: %v", err))
|
|
return
|
|
}
|
|
|
|
// Загружаем как новый архив
|
|
newDeploySpec, err := s.buildDeploySpec(ctx, archiveBytes)
|
|
if err != nil {
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("upload clone archive: %v", err))
|
|
return
|
|
}
|
|
|
|
// Параметры для нового пакета (берём spec из оригинала)
|
|
now := time.Now().UTC()
|
|
newPkgName := req.NewName + "-" + now.Format("20060102150405")
|
|
|
|
// Определяем environment из исходной функции
|
|
envName, _, _ := unstructured.NestedString(srcFn.Object, "spec", "environment", "name")
|
|
|
|
// Собираем spec пакета (аналогично оригиналу, но с новыми байтами)
|
|
// Если у оригинала есть source (Go) — копируем source spec
|
|
srcSourceSpec, _, _ := unstructured.NestedMap(srcPkg.Object, "spec", "source")
|
|
hasBuildCmd, _, _ := unstructured.NestedString(srcPkg.Object, "spec", "buildcommand")
|
|
|
|
var newPkgSpec map[string]any
|
|
if hasBuildCmd != "" {
|
|
// Go: source package
|
|
newPkgSpec = map[string]any{
|
|
"source": newDeploySpec, // перезаливаем source
|
|
"deployment": map[string]any{},
|
|
"environment": map[string]any{"name": envName, "namespace": ns},
|
|
"buildcommand": hasBuildCmd,
|
|
}
|
|
_ = srcSourceSpec
|
|
} else {
|
|
newPkgSpec = map[string]any{
|
|
"deployment": newDeploySpec,
|
|
"environment": map[string]any{"name": envName, "namespace": ns},
|
|
"source": map[string]any{},
|
|
}
|
|
}
|
|
|
|
newPkg := &unstructured.Unstructured{Object: map[string]any{
|
|
"apiVersion": "fission.io/v1",
|
|
"kind": "Package",
|
|
"metadata": map[string]any{"name": newPkgName, "namespace": ns},
|
|
"spec": newPkgSpec,
|
|
}}
|
|
if _, err := s.dyn.Resource(fission.PackageGVR).Namespace(ns).Create(ctx, newPkg, metav1.CreateOptions{}); err != nil {
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create clone package: %v", err))
|
|
return
|
|
}
|
|
|
|
// Копируем аннотации из исходной функции
|
|
srcAnnotations := srcFn.GetAnnotations()
|
|
newAnnotations := map[string]any{
|
|
functionCreatedAtAnnotation: now.Format(time.RFC3339),
|
|
functionUpdatedAtAnnotation: now.Format(time.RFC3339),
|
|
}
|
|
for _, k := range []string{
|
|
"fission-console/language",
|
|
fissionSourceTypeAnnotation,
|
|
"fission-console/env-vars",
|
|
"fission-console/archive-filename",
|
|
} {
|
|
if v, ok := srcAnnotations[k]; ok && v != "" {
|
|
newAnnotations[k] = v
|
|
}
|
|
}
|
|
newAnnotations["fission-console/cloned-from"] = srcName
|
|
|
|
// Копируем entrypoint
|
|
entrypoint, _, _ := unstructured.NestedString(srcFn.Object, "spec", "package", "functionName")
|
|
timeout, _, _ := unstructured.NestedInt64(srcFn.Object, "spec", "functionTimeout")
|
|
if timeout == 0 {
|
|
timeout = 60
|
|
}
|
|
|
|
// Копируем InvokeStrategy и podspec
|
|
invokeStrategy, _, _ := unstructured.NestedMap(srcFn.Object, "spec", "InvokeStrategy")
|
|
if invokeStrategy == nil {
|
|
invokeStrategy = map[string]any{
|
|
"ExecutionStrategy": map[string]any{"ExecutorType": "poolmgr"},
|
|
"StrategyType": "execution",
|
|
}
|
|
}
|
|
podspec, _, _ := unstructured.NestedMap(srcFn.Object, "spec", "podspec")
|
|
|
|
newFnSpec := map[string]any{
|
|
"environment": map[string]any{"name": envName, "namespace": ns},
|
|
"functionTimeout": timeout,
|
|
"InvokeStrategy": invokeStrategy,
|
|
"package": map[string]any{
|
|
"packageref": map[string]any{"name": newPkgName, "namespace": ns},
|
|
"functionName": entrypoint,
|
|
},
|
|
}
|
|
if len(podspec) > 0 {
|
|
newFnSpec["podspec"] = podspec
|
|
}
|
|
|
|
newFn := &unstructured.Unstructured{Object: map[string]any{
|
|
"apiVersion": "fission.io/v1",
|
|
"kind": "Function",
|
|
"metadata": map[string]any{
|
|
"name": req.NewName,
|
|
"namespace": ns,
|
|
"annotations": newAnnotations,
|
|
},
|
|
"spec": newFnSpec,
|
|
}}
|
|
if _, err := s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Create(ctx, newFn, metav1.CreateOptions{}); err != nil {
|
|
_ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, newPkgName, metav1.DeleteOptions{})
|
|
if apierrors.IsAlreadyExists(err) {
|
|
writeJSONError(w, http.StatusConflict, fmt.Sprintf("function %q already exists", req.NewName))
|
|
return
|
|
}
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create clone function: %v", err))
|
|
return
|
|
}
|
|
|
|
// Маршрут для нового триггера
|
|
if req.Route == "" {
|
|
nsShort := ns
|
|
if len(nsShort) > 12 {
|
|
nsShort = nsShort[len(nsShort)-12:]
|
|
}
|
|
req.Route = "/" + nsShort + "/" + req.NewName
|
|
}
|
|
if !strings.HasPrefix(req.Route, "/") {
|
|
req.Route = "/" + req.Route
|
|
}
|
|
triggerName := req.NewName + "-route"
|
|
|
|
// Определяем методы из существующего триггера оригинала
|
|
methods := s.getTriggerMethods(ctx, ns, srcName)
|
|
if len(methods) == 0 {
|
|
methods = []any{"GET"}
|
|
}
|
|
|
|
newTrigger := &unstructured.Unstructured{Object: map[string]any{
|
|
"apiVersion": "fission.io/v1",
|
|
"kind": "HTTPTrigger",
|
|
"metadata": map[string]any{"name": triggerName, "namespace": ns},
|
|
"spec": map[string]any{
|
|
"relativeurl": req.Route,
|
|
"methods": methods,
|
|
"createingress": true,
|
|
"functionref": map[string]any{"type": "name", "name": req.NewName},
|
|
},
|
|
}}
|
|
if _, err := s.dyn.Resource(fission.HTTPTrigGVR).Namespace(ns).Create(ctx, newTrigger, metav1.CreateOptions{}); err != nil {
|
|
_ = s.dyn.Resource(fission.FunctionGVR).Namespace(ns).Delete(ctx, req.NewName, metav1.DeleteOptions{})
|
|
_ = s.dyn.Resource(fission.PackageGVR).Namespace(ns).Delete(ctx, newPkgName, metav1.DeleteOptions{})
|
|
writeJSONError(w, http.StatusBadGateway, fmt.Sprintf("create httptrigger: %v", err))
|
|
return
|
|
}
|
|
|
|
s.billing.RecordInvocation(billing.Invocation{
|
|
Namespace: ns,
|
|
FunctionName: req.NewName,
|
|
TriggerType: billing.TriggerEvent,
|
|
StartedAt: now,
|
|
StatusCode: http.StatusCreated,
|
|
RecordedBy: "console",
|
|
EventType: "clone",
|
|
})
|
|
|
|
writeAnyJSON(w, http.StatusCreated, map[string]any{
|
|
"name": req.NewName,
|
|
"cloned_from": srcName,
|
|
"package": newPkgName,
|
|
"route": req.Route,
|
|
})
|
|
}
|
|
|
|
// downloadPackageBytes извлекает байты архива из Package CRD.
|
|
// Поддерживает type:url (скачивает из storagesvc) и type:literal (base64).
|
|
func (s *Server) downloadPackageBytes(ctx context.Context, pkg *unstructured.Unstructured) ([]byte, error) {
|
|
// Пробуем deployment сначала, потом source (для Go)
|
|
for _, field := range [][]string{{"spec", "deployment"}, {"spec", "source"}} {
|
|
spec, _, _ := unstructured.NestedMap(pkg.Object, field...)
|
|
if len(spec) == 0 {
|
|
continue
|
|
}
|
|
archiveType, _ := spec["type"].(string)
|
|
switch archiveType {
|
|
case "url":
|
|
archiveURL, _ := spec["url"].(string)
|
|
if archiveURL == "" {
|
|
continue
|
|
}
|
|
return s.downloadFromStoragesvc(ctx, archiveURL)
|
|
case "literal":
|
|
lit, _ := spec["literal"].(string)
|
|
if lit == "" {
|
|
continue
|
|
}
|
|
return base64.StdEncoding.DecodeString(lit)
|
|
}
|
|
}
|
|
// Последний шанс: если функция Python с простым кодом — возвращаем заглушку
|
|
return nil, fmt.Errorf("no downloadable archive found in package (empty deployment and source spec)")
|
|
}
|
|
|
|
// downloadFromStoragesvc скачивает архив по URL из storagesvc.
|
|
func (s *Server) downloadFromStoragesvc(ctx context.Context, archiveURL string) ([]byte, error) {
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, archiveURL, nil)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("build download request: %w", err)
|
|
}
|
|
resp, err := s.http.Do(req)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("download archive: %w", err)
|
|
}
|
|
defer resp.Body.Close()
|
|
if resp.StatusCode != http.StatusOK {
|
|
return nil, fmt.Errorf("download archive status %d", resp.StatusCode)
|
|
}
|
|
data, err := io.ReadAll(resp.Body)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("read archive body: %w", err)
|
|
}
|
|
return data, nil
|
|
}
|
|
|
|
// getTriggerMethods возвращает методы HTTP из триггера функции (для копирования в клон).
|
|
func (s *Server) getTriggerMethods(ctx context.Context, ns, fnName string) []any {
|
|
triggers, err := s.dyn.Resource(fission.HTTPTrigGVR).Namespace(ns).List(ctx, metav1.ListOptions{})
|
|
if err != nil {
|
|
return nil
|
|
}
|
|
for _, t := range triggers.Items {
|
|
ref, _, _ := unstructured.NestedString(t.Object, "spec", "functionref", "name")
|
|
if ref != fnName {
|
|
continue
|
|
}
|
|
methods, _, _ := unstructured.NestedSlice(t.Object, "spec", "methods")
|
|
if len(methods) > 0 {
|
|
return methods
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Ссылка на runtime.DefaultEntrypoint для возможного использования в будущем
|
|
var _ = runtime.DefaultEntrypoint
|