diff --git a/Makefile b/Makefile index 488287a9..f45df4f8 100644 --- a/Makefile +++ b/Makefile @@ -10,7 +10,14 @@ GO_BIN ?= $(shell go env GOPATH)/bin DRUID_K8S_NAMESPACE ?= druid DRUID_WORKER_CALLBACK_LISTEN ?= 0.0.0.0:8083 DRUID_WORKER_CALLBACK_URL ?= http://host.k3d.internal:8083 -DRUID_WATCH_ARGS ?= daemon --runtime kubernetes --listen 127.0.0.1:8081 --public-listen 127.0.0.1:8082 --worker-callback-listen $(DRUID_WORKER_CALLBACK_LISTEN) --worker-callback-url $(DRUID_WORKER_CALLBACK_URL) --unsafe-allow-unauthenticated-management --unsafe-allow-unauthenticated-public --k8s-namespace $(DRUID_K8S_NAMESPACE) --k8s-pull-image $(DRUID_K8S_PULL_IMAGE) +DRUID_WORKER_DAEMON_URL ?= http://host.k3d.internal:8081 +DRUID_WATCH_ARGS ?= daemon --runtime kubernetes --listen 127.0.0.1:8081 --public-listen 127.0.0.1:8082 --worker-callback-listen $(DRUID_WORKER_CALLBACK_LISTEN) --worker-callback-url $(DRUID_WORKER_CALLBACK_URL) --worker-daemon-url $(DRUID_WORKER_DAEMON_URL) --unsafe-allow-unauthenticated-management --unsafe-allow-unauthenticated-public --k8s-namespace $(DRUID_K8S_NAMESPACE) --k8s-pull-image $(DRUID_K8S_PULL_IMAGE) +export DRUID_K8S_UI_S3_BUCKET ?= druid-ui +export DRUID_K8S_UI_S3_PUBLIC_BASE_URL ?= http://127.0.0.1:9000/druid-ui +export DRUID_K8S_UI_S3_REGION ?= us-east-1 +export DRUID_K8S_UI_S3_ENDPOINT ?= http://127.0.0.1:9000 +export DRUID_K8S_UI_S3_ACCESS_KEY ?= druid +export DRUID_K8S_UI_S3_SECRET_KEY ?= druidpassword generate-api: ## Generate API types from OpenAPI spec @echo "Generating API types from OpenAPI spec..." diff --git a/apps/druid/adapters/cli/client/dev.go b/apps/druid/adapters/cli/client/dev.go index e1d8a83d..5bdd670c 100644 --- a/apps/druid/adapters/cli/client/dev.go +++ b/apps/druid/adapters/cli/client/dev.go @@ -2,6 +2,7 @@ package client import ( "context" + "crypto/subtle" "encoding/json" "fmt" "mime" @@ -130,7 +131,7 @@ func runDevServer() error { if devDaemonToken == "" { devDaemonToken = os.Getenv("DRUID_INTERNAL_TOKEN") } - auth := devAuth{runtimeID: devRuntimeID, ownerID: devOwnerID} + auth := devAuth{runtimeID: devRuntimeID, ownerID: devOwnerID, internalToken: devDaemonToken} if devAuthJWKSURL != "" { auth.user, err = coreservices.NewAuthorizer([]string{devAuthJWKSURL}, "") if err != nil { @@ -161,10 +162,11 @@ func runDevServer() error { } type devAuth struct { - user ports.AuthorizerServiceInterface - runtime ports.AuthorizerServiceInterface - runtimeID string - ownerID string + user ports.AuthorizerServiceInterface + runtime ports.AuthorizerServiceInterface + runtimeID string + ownerID string + internalToken string } func newDevApp(root string, broadcast *domain.BroadcastChannel, queue *devTriggerQueue, authOpt ...devAuth) *fiber.App { @@ -190,6 +192,7 @@ func newDevApp(root string, broadcast *domain.BroadcastChannel, queue *devTrigge }) app.Use(server.authMiddleware) devapi.RegisterHandlers(app, server) + app.Get("/internal/v1/ui/*", server.GetInternalUIPackage) webdavHandler := adaptor.HTTPHandler(&webdav.Handler{ Prefix: "/webdav", FileSystem: webdav.Dir(root), @@ -223,6 +226,13 @@ func (s devServer) authMiddleware(c *fiber.Ctx) error { if c.Path() == "/health" || c.Method() == fiber.MethodOptions { return c.Next() } + if strings.HasPrefix(c.Path(), "/internal/v1/ui/") { + token := strings.TrimPrefix(c.Get("Authorization"), "Bearer ") + if s.auth.internalToken != "" && subtle.ConstantTimeCompare([]byte(token), []byte(s.auth.internalToken)) != 1 { + return fiber.NewError(fiber.StatusUnauthorized, "invalid internal token") + } + return c.Next() + } if s.auth.user == nil && s.auth.runtime == nil { return c.Next() } @@ -252,6 +262,15 @@ func (s devServer) authMiddleware(c *fiber.Ctx) error { return c.Next() } +func (s devServer) GetInternalUIPackage(c *fiber.Ctx) error { + path := filepath.ToSlash(filepath.Clean(strings.TrimPrefix(c.Params("*"), "/"))) + if path == "." || strings.HasPrefix(path, "../") || filepath.Ext(path) != ".wasm" || + (!strings.HasPrefix(path, "private/") && !strings.HasPrefix(path, "public/")) { + return fiber.ErrNotFound + } + return s.sendFile(c, filepath.ToSlash(filepath.Join(domain.RuntimeDataDir, path))) +} + func (s devServer) GetFile(c *fiber.Ctx, params devapi.GetFileParams) error { return s.sendFile(c, params.Path) } diff --git a/apps/druid/adapters/cli/client/dev_test.go b/apps/druid/adapters/cli/client/dev_test.go index 6f722840..96264caa 100644 --- a/apps/druid/adapters/cli/client/dev_test.go +++ b/apps/druid/adapters/cli/client/dev_test.go @@ -195,6 +195,36 @@ func TestDevServerFileAuth(t *testing.T) { } } +func TestDevServerInternalUIPackageRequiresDaemonToken(t *testing.T) { + root := t.TempDir() + packagePath := filepath.Join(root, "data", "private", "dist") + if err := os.MkdirAll(packagePath, 0o755); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(packagePath, "app.wasm"), []byte("wasm"), 0o644); err != nil { + t.Fatal(err) + } + app := newDevApp(root, domain.NewHub(), &devTriggerQueue{}, devAuth{internalToken: "internal-token"}) + + response, err := app.Test(httptest.NewRequest(http.MethodGet, "/internal/v1/ui/private/dist/app.wasm", nil)) + if err != nil { + t.Fatal(err) + } + if response.StatusCode != http.StatusUnauthorized { + t.Fatalf("status = %d, want %d", response.StatusCode, http.StatusUnauthorized) + } + + request := httptest.NewRequest(http.MethodGet, "/internal/v1/ui/private/dist/app.wasm", nil) + request.Header.Set("Authorization", "Bearer internal-token") + response, err = app.Test(request) + if err != nil { + t.Fatal(err) + } + if response.StatusCode != http.StatusOK { + t.Fatalf("status = %d, want %d", response.StatusCode, http.StatusOK) + } +} + type devTestAuth struct{} func (devTestAuth) CheckHeader(c *fiber.Ctx) (*ports.AuthContext, error) { diff --git a/apps/druid/adapters/cli/daemon.go b/apps/druid/adapters/cli/daemon.go index 7369d314..a457e242 100644 --- a/apps/druid/adapters/cli/daemon.go +++ b/apps/druid/adapters/cli/daemon.go @@ -34,7 +34,9 @@ var k8sUIS3PublicBaseURL string var k8sUIS3Region string var k8sUIS3Endpoint string var k8sUIS3Prefix string -var k8sUIS3Secret string +var k8sUIS3AccessKey string +var k8sUIS3SecretKey string +var k8sUIS3SessionToken string var k8sKubeconfig string var runtimeListen string var runtimePublicListen string @@ -95,7 +97,9 @@ func init() { DaemonCommand.Flags().StringVar(&k8sUIS3Region, "k8s-ui-s3-region", "", "S3 region for published UI packages (default: DRUID_K8S_UI_S3_REGION)") DaemonCommand.Flags().StringVar(&k8sUIS3Endpoint, "k8s-ui-s3-endpoint", "", "Optional S3-compatible endpoint for UI packages (default: DRUID_K8S_UI_S3_ENDPOINT)") DaemonCommand.Flags().StringVar(&k8sUIS3Prefix, "k8s-ui-s3-prefix", "", "Optional S3 key prefix for UI packages (default: DRUID_K8S_UI_S3_PREFIX)") - DaemonCommand.Flags().StringVar(&k8sUIS3Secret, "k8s-ui-s3-credentials-secret", "", "Kubernetes secret with AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY (default: DRUID_K8S_UI_S3_CREDENTIALS_SECRET)") + DaemonCommand.Flags().StringVar(&k8sUIS3AccessKey, "k8s-ui-s3-access-key", "", "S3 access key for published UI packages (default: DRUID_K8S_UI_S3_ACCESS_KEY)") + DaemonCommand.Flags().StringVar(&k8sUIS3SecretKey, "k8s-ui-s3-secret-key", "", "S3 secret key for published UI packages (default: DRUID_K8S_UI_S3_SECRET_KEY)") + DaemonCommand.Flags().StringVar(&k8sUIS3SessionToken, "k8s-ui-s3-session-token", "", "Optional S3 session token for published UI packages (default: DRUID_K8S_UI_S3_SESSION_TOKEN)") DaemonCommand.Flags().StringVar(&k8sKubeconfig, "k8s-kubeconfig", "", "Kubernetes kubeconfig path for out-of-cluster runtime access (default: DRUID_K8S_KUBECONFIG, KUBECONFIG, or ~/.kube/config)") } @@ -115,7 +119,10 @@ func runRuntimeDaemon() error { UIS3Region: k8sUIS3Region, UIS3Endpoint: k8sUIS3Endpoint, UIS3Prefix: k8sUIS3Prefix, - UIS3Secret: k8sUIS3Secret, + UIS3AccessKey: k8sUIS3AccessKey, + UIS3SecretKey: k8sUIS3SecretKey, + UIS3SessionToken: k8sUIS3SessionToken, + InternalToken: runtimeInternalToken, } dockerConfig := runtimedocker.Config{WorkerImage: dockerWorkerImage, Storage: dockerStorage, BindRoot: dockerBindRoot, VolumePrefix: dockerVolumePrefix, UIBind: dockerUIBind, UIPublicURL: dockerUIPublicURL} logManager := services.NewLogManager() diff --git a/apps/druid/adapters/cli/ui.go b/apps/druid/adapters/cli/ui.go index e4353a0f..71dadd68 100644 --- a/apps/druid/adapters/cli/ui.go +++ b/apps/druid/adapters/cli/ui.go @@ -1,20 +1,14 @@ package cli import ( - "bytes" "context" - "crypto/sha256" - "encoding/hex" "fmt" "net/http" "os" - "path" "path/filepath" "strings" - "github.com/aws/aws-sdk-go-v2/aws" - awscfg "github.com/aws/aws-sdk-go-v2/config" - "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/highcard-dev/daemon/internal/uipackage" "github.com/spf13/cobra" ) @@ -99,30 +93,9 @@ func publishUIPackageToS3(ctx context.Context) (string, error) { } return "", err } - sum := sha256.Sum256(data) - hash := hex.EncodeToString(sum[:]) - cfg, err := awscfg.LoadDefaultConfig(ctx, awscfg.WithRegion(uiPublishRegion)) - if err != nil { - return "", err - } - client := s3.NewFromConfig(cfg, func(options *s3.Options) { - if uiPublishEndpoint != "" { - options.BaseEndpoint = aws.String(uiPublishEndpoint) - options.UsePathStyle = true - } + return uipackage.Upload(ctx, data, uipackage.S3Config{ + Bucket: uiPublishBucket, Region: uiPublishRegion, Endpoint: uiPublishEndpoint, KeyPrefix: uiPublishKeyPrefix, }) - key := path.Join(strings.Trim(uiPublishKeyPrefix, "/"), hash, "app.wasm") - _, err = client.PutObject(ctx, &s3.PutObjectInput{ - Bucket: aws.String(uiPublishBucket), - Key: aws.String(key), - Body: bytes.NewReader(data), - ContentType: aws.String("application/wasm"), - CacheControl: aws.String("public, max-age=31536000, immutable"), - }) - if err != nil { - return "", err - } - return hash, nil } func cleanUIPackageSource(source string) (string, error) { diff --git a/internal/runtime/docker/ui.go b/internal/runtime/docker/ui.go index c29cf900..5d13da32 100644 --- a/internal/runtime/docker/ui.go +++ b/internal/runtime/docker/ui.go @@ -12,6 +12,7 @@ import ( "github.com/docker/docker/api/types/mount" "github.com/docker/docker/pkg/stdcopy" "github.com/docker/go-connections/nat" + "github.com/highcard-dev/daemon/internal/core/domain" "github.com/highcard-dev/daemon/internal/core/ports" ) @@ -87,7 +88,7 @@ func (b *Backend) ensureUIPackageServer(ctx context.Context) error { } func (b *Backend) copyUIPackage(ctx context.Context, action ports.RuntimeUIPackageAction) (string, error) { - rootMount, err := DockerMount(action.RootRef, "/scroll", true, "") + rootMount, err := DockerMount(action.RootRef, "/scroll", true, domain.RuntimeDataDir) if err != nil { return "", err } diff --git a/internal/runtime/kubernetes/backend.go b/internal/runtime/kubernetes/backend.go index f5f6e20c..bbeae470 100644 --- a/internal/runtime/kubernetes/backend.go +++ b/internal/runtime/kubernetes/backend.go @@ -3,6 +3,7 @@ package kubernetes import ( "context" "fmt" + "net/http" "sync" "time" @@ -12,19 +13,23 @@ import ( "k8s.io/client-go/tools/clientcmd" "github.com/highcard-dev/daemon/internal/core/ports" + "github.com/highcard-dev/daemon/internal/uipackage" "github.com/highcard-dev/daemon/internal/utils/logger" "go.uber.org/zap" ) type Backend struct { - client k8sclient.Interface - restConfig *rest.Config - consoleManager ports.ConsoleManagerInterface - config Config - statsReader nodeStatsReader - jobLogRunner func(context.Context, *batchv1.Job) ([]byte, error) - jobExitMu sync.Mutex - jobExits map[string]recentJobExit + client k8sclient.Interface + restConfig *rest.Config + httpClient *http.Client + consoleManager ports.ConsoleManagerInterface + config Config + statsReader nodeStatsReader + jobLogRunner func(context.Context, *batchv1.Job) ([]byte, error) + uiPackageFetcher func(context.Context, string, string, string, string) ([]byte, error) + uiPackageUploader func(context.Context, []byte, uipackage.S3Config) (string, error) + jobExitMu sync.Mutex + jobExits map[string]recentJobExit } type recentJobExit struct { @@ -56,10 +61,15 @@ func New(config Config, consoleManager ports.ConsoleManagerInterface) (*Backend, if _, err := client.Discovery().ServerVersion(); err != nil { return nil, fmt.Errorf("kubernetes API unavailable: %w", err) } + httpClient, err := rest.HTTPClientFor(restConfig) + if err != nil { + return nil, fmt.Errorf("kubernetes HTTP client unavailable: %w", err) + } logger.Log().Info("Using Kubernetes backend settings", zap.String("source", source), zap.String("namespace", config.Namespace)) backend := &Backend{ client: client, restConfig: restConfig, + httpClient: httpClient, consoleManager: consoleManager, config: config, jobExits: make(map[string]recentJobExit), diff --git a/internal/runtime/kubernetes/config.go b/internal/runtime/kubernetes/config.go index 50c350ed..05725f36 100644 --- a/internal/runtime/kubernetes/config.go +++ b/internal/runtime/kubernetes/config.go @@ -23,7 +23,10 @@ type Config struct { UIS3Region string UIS3Endpoint string UIS3Prefix string - UIS3Secret string + UIS3AccessKey string + UIS3SecretKey string + UIS3SessionToken string + InternalToken string } func (c Config) WithDefaults() Config { @@ -66,8 +69,17 @@ func (c Config) WithDefaults() Config { if c.UIS3Prefix == "" { c.UIS3Prefix = os.Getenv("DRUID_K8S_UI_S3_PREFIX") } - if c.UIS3Secret == "" { - c.UIS3Secret = os.Getenv("DRUID_K8S_UI_S3_CREDENTIALS_SECRET") + if c.UIS3AccessKey == "" { + c.UIS3AccessKey = os.Getenv("DRUID_K8S_UI_S3_ACCESS_KEY") + } + if c.UIS3SecretKey == "" { + c.UIS3SecretKey = os.Getenv("DRUID_K8S_UI_S3_SECRET_KEY") + } + if c.UIS3SessionToken == "" { + c.UIS3SessionToken = os.Getenv("DRUID_K8S_UI_S3_SESSION_TOKEN") + } + if c.InternalToken == "" { + c.InternalToken = os.Getenv("DRUID_INTERNAL_TOKEN") } return c } @@ -95,8 +107,8 @@ func (c Config) ValidateForUIPublishing() error { if c.PullImage == "" { return fmt.Errorf("kubernetes pull image is required for UI publishing; set --k8s-pull-image or DRUID_K8S_PULL_IMAGE") } - if c.UIS3Bucket == "" || c.UIS3PublicBaseURL == "" || c.UIS3Region == "" || c.UIS3Secret == "" { - return fmt.Errorf("kubernetes UI publishing requires DRUID_K8S_UI_S3_BUCKET, DRUID_K8S_UI_S3_PUBLIC_BASE_URL, DRUID_K8S_UI_S3_REGION, and DRUID_K8S_UI_S3_CREDENTIALS_SECRET") + if c.UIS3Bucket == "" || c.UIS3PublicBaseURL == "" || c.UIS3Region == "" || c.UIS3AccessKey == "" || c.UIS3SecretKey == "" { + return fmt.Errorf("kubernetes UI publishing requires S3 bucket, public URL, region, access key, and secret key configuration") } return nil } diff --git a/internal/runtime/kubernetes/resources.go b/internal/runtime/kubernetes/resources.go index 575efa79..06a06401 100644 --- a/internal/runtime/kubernetes/resources.go +++ b/internal/runtime/kubernetes/resources.go @@ -256,6 +256,7 @@ func devStatefulSetSpec(namespace string, root string, pvc string, image string, labels := baseLabels(pvc) labels[labelProcedure] = "dev" replicas := int32(1) + runAsRoot := int64(0) args := []string{"dev", "--root", action.MountPath, "--listen", action.Listen, "--runtime-id", action.RuntimeID, "--daemon-url", action.DaemonURL} if action.DaemonToken != "" { args = append(args, "--daemon-token", action.DaemonToken) @@ -282,6 +283,7 @@ func devStatefulSetSpec(namespace string, root string, pvc string, image string, Command: []string{"druid"}, Args: args, ImagePullPolicy: corev1.PullIfNotPresent, + SecurityContext: &corev1.SecurityContext{RunAsUser: &runAsRoot}, Ports: []corev1.ContainerPort{{Name: "webdav", ContainerPort: 8084}}, VolumeMounts: []corev1.VolumeMount{{Name: "data", MountPath: action.MountPath}}, }}, diff --git a/internal/runtime/kubernetes/resources_test.go b/internal/runtime/kubernetes/resources_test.go index af7b3512..98a556f1 100644 --- a/internal/runtime/kubernetes/resources_test.go +++ b/internal/runtime/kubernetes/resources_test.go @@ -21,6 +21,7 @@ import ( "github.com/highcard-dev/daemon/internal/core/domain" "github.com/highcard-dev/daemon/internal/core/ports" coreservices "github.com/highcard-dev/daemon/internal/core/services" + "github.com/highcard-dev/daemon/internal/uipackage" ) func setStatsReader(backend *Backend, nodeName string, summary *nodeStatsSummary, err error) { @@ -122,6 +123,58 @@ func TestProcedureJobSpecBuildsDeterministicMountsAndLabels(t *testing.T) { } } +func TestDevStatefulSetMountsRuntimeRootAsScrollRoot(t *testing.T) { + root := ref("druid", "druid-ui-data") + statefulSet := devStatefulSetSpec("druid", root, "druid-ui-data", "druid:local", ports.RuntimeDevAction{ + RuntimeID: "ui", RootRef: root, MountPath: "/scroll", Listen: ":8084", + }, "registry-secret") + mount := statefulSet.Spec.Template.Spec.Containers[0].VolumeMounts[0] + if mount.MountPath != "/scroll" || mount.SubPath != "" { + t.Fatalf("mount = %#v, want the runtime root at /scroll", mount) + } + if user := statefulSet.Spec.Template.Spec.Containers[0].SecurityContext.RunAsUser; user == nil || *user != 0 { + t.Fatalf("dev server must run as root to edit runtime files, got %#v", user) + } +} + +func TestPublishUIPackageFetchesFromDevServiceAndUploadsFromDaemon(t *testing.T) { + root := ref("druid", "druid-ui-data") + client := fake.NewSimpleClientset(&appsv1.StatefulSet{ObjectMeta: metav1.ObjectMeta{ + Name: devStatefulSetName(root), Namespace: "druid", + }}) + backend := NewWithClient(Config{ + PullImage: "druid:local", + UIS3Bucket: "druid-ui", + UIS3PublicBaseURL: "http://127.0.0.1:9000/druid-ui", + UIS3Region: "us-east-1", + UIS3AccessKey: "druid", + UIS3SecretKey: "druidpassword", + InternalToken: "internal-token", + }, nil, client) + backend.uiPackageFetcher = func(_ context.Context, namespace string, service string, source string, token string) ([]byte, error) { + if namespace != "druid" || service != serviceName(root, "dev", "webdav") || source != "private/dist/app.wasm" || token != "internal-token" { + t.Fatalf("fetch = namespace=%q service=%q source=%q token=%q", namespace, service, source, token) + } + return []byte("wasm"), nil + } + backend.uiPackageUploader = func(_ context.Context, content []byte, config uipackage.S3Config) (string, error) { + if string(content) != "wasm" || config.Bucket != "druid-ui" || config.AccessKey != "druid" || config.SecretKey != "druidpassword" { + t.Fatalf("upload content=%q config=%#v", content, config) + } + return "content-hash", nil + } + + result, err := backend.PublishUIPackage(context.Background(), ports.RuntimeUIPackageAction{ + RuntimeID: "ui", RootRef: root, Scope: domain.RuntimeUIPackageScopePrivate, SourcePath: "private/dist/app.wasm", + }) + if err != nil { + t.Fatal(err) + } + if result.URL != "http://127.0.0.1:9000/druid-ui/druid/ui/private/content-hash/app.wasm" { + t.Fatalf("url = %q", result.URL) + } +} + func TestProcedureJobSpecUsesProvidedRuntimeEnv(t *testing.T) { procedure := &domain.Procedure{ Image: "alpine:3.20", diff --git a/internal/runtime/kubernetes/ui.go b/internal/runtime/kubernetes/ui.go index 8baab341..54f74354 100644 --- a/internal/runtime/kubernetes/ui.go +++ b/internal/runtime/kubernetes/ui.go @@ -3,13 +3,17 @@ package kubernetes import ( "context" "fmt" + "io" + "net/http" + "net/url" "path" "strings" - batchv1 "k8s.io/api/batch/v1" - corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "github.com/highcard-dev/daemon/internal/core/ports" + "github.com/highcard-dev/daemon/internal/uipackage" ) func (b *Backend) PublishUIPackage(ctx context.Context, action ports.RuntimeUIPackageAction) (ports.RuntimeUIPackageResult, error) { @@ -20,18 +24,32 @@ func (b *Backend) PublishUIPackage(ctx context.Context, action ports.RuntimeUIPa if err != nil { return ports.RuntimeUIPackageResult{}, err } - keyPrefix := strings.Trim(strings.Trim(b.config.UIS3Prefix, "/")+"/"+namespace+"/"+action.RuntimeID+"/"+string(action.Scope), "/") - job := uiPublishJobSpec(namespace, jobName("ui-publish", action.RootRef, string(action.Scope)), pvc, b.config.PullImage, action, b.config.RegistrySecret, b.config.UIS3Secret, b.config.UIS3Bucket, b.config.UIS3Region, b.config.UIS3Endpoint, keyPrefix) - logs, err := b.runJobAndLogs(ctx, job) + temporaryDevServer, err := b.ensureUIPackageDevServer(ctx, namespace, pvc, action) if err != nil { return ports.RuntimeUIPackageResult{}, err } - hash := strings.TrimSpace(string(logs)) - if idx := strings.LastIndex(hash, "\n"); idx >= 0 { - hash = strings.TrimSpace(hash[idx+1:]) + if temporaryDevServer { + defer b.StopDev(context.Background(), action.RootRef) + } + content, err := b.fetchUIPackage(ctx, namespace, action.RootRef, action.SourcePath) + if err != nil { + return ports.RuntimeUIPackageResult{}, err + } + keyPrefix := strings.Trim(strings.Trim(b.config.UIS3Prefix, "/")+"/"+namespace+"/"+action.RuntimeID+"/"+string(action.Scope), "/") + hash, err := b.uploadUIPackage(ctx, content, uipackage.S3Config{ + Bucket: b.config.UIS3Bucket, + Region: b.config.UIS3Region, + Endpoint: b.config.UIS3Endpoint, + KeyPrefix: keyPrefix, + AccessKey: b.config.UIS3AccessKey, + SecretKey: b.config.UIS3SecretKey, + SessionToken: b.config.UIS3SessionToken, + }) + if err != nil { + return ports.RuntimeUIPackageResult{}, err } if hash == "" { - return ports.RuntimeUIPackageResult{}, fmt.Errorf("ui publish job did not return a content hash") + return ports.RuntimeUIPackageResult{}, fmt.Errorf("ui package upload did not return a content hash") } key := path.Join(keyPrefix, hash, "app.wasm") return ports.RuntimeUIPackageResult{ @@ -41,50 +59,70 @@ func (b *Backend) PublishUIPackage(ctx context.Context, action ports.RuntimeUIPa }, nil } -func boolPtr(value bool) *bool { - return &value +func (b *Backend) ensureUIPackageDevServer(ctx context.Context, namespace string, pvc string, action ports.RuntimeUIPackageAction) (bool, error) { + name := devStatefulSetName(action.RootRef) + _, err := b.client.AppsV1().StatefulSets(namespace).Get(ctx, name, metav1.GetOptions{}) + if err == nil { + return false, nil + } + if !apierrors.IsNotFound(err) { + return false, err + } + devAction := ports.RuntimeDevAction{ + RuntimeID: action.RuntimeID, + RootRef: action.RootRef, + MountPath: "/scroll", + Listen: ":8084", + DaemonToken: b.config.InternalToken, + } + server := devStatefulSetSpec(namespace, action.RootRef, pvc, b.config.PullImage, devAction, b.config.RegistrySecret) + if _, err := b.client.AppsV1().StatefulSets(namespace).Create(ctx, server, metav1.CreateOptions{}); err != nil { + return false, err + } + if err := b.reconcileService(ctx, devServiceSpec(namespace, action.RootRef, pvc)); err != nil { + _ = b.StopDev(context.Background(), action.RootRef) + return false, err + } + if err := b.waitForStatefulSet(ctx, namespace, name); err != nil { + _ = b.StopDev(context.Background(), action.RootRef) + return false, err + } + return true, nil } -func uiPublishJobSpec(namespace string, jobName string, pvc string, image string, action ports.RuntimeUIPackageAction, imagePullSecret string, s3Secret string, bucket string, region string, endpoint string, keyPrefix string) *batchv1.Job { - command := []string{ - "druid", "ui", "publish-s3", - "--root", "/scroll", - "--source", action.SourcePath, - "--bucket", bucket, - "--region", region, - "--key-prefix", keyPrefix, - } - if endpoint != "" { - command = append(command, "--endpoint", endpoint) - } - job := helperJobSpec(namespace, jobName, pvc, image, command, imagePullSecret, map[string]string{ - labelComponent: "ui-publish", - labelRuntimeID: dnsLabel(action.RuntimeID), - }) - container := &job.Spec.Template.Spec.Containers[0] - container.Env = append(container.Env, - corev1.EnvVar{ - Name: "AWS_ACCESS_KEY_ID", - ValueFrom: &corev1.EnvVarSource{SecretKeyRef: &corev1.SecretKeySelector{ - LocalObjectReference: corev1.LocalObjectReference{Name: s3Secret}, - Key: "AWS_ACCESS_KEY_ID", - }}, - }, - corev1.EnvVar{ - Name: "AWS_SECRET_ACCESS_KEY", - ValueFrom: &corev1.EnvVarSource{SecretKeyRef: &corev1.SecretKeySelector{ - LocalObjectReference: corev1.LocalObjectReference{Name: s3Secret}, - Key: "AWS_SECRET_ACCESS_KEY", - }}, - }, - corev1.EnvVar{ - Name: "AWS_SESSION_TOKEN", - ValueFrom: &corev1.EnvVarSource{SecretKeyRef: &corev1.SecretKeySelector{ - LocalObjectReference: corev1.LocalObjectReference{Name: s3Secret}, - Key: "AWS_SESSION_TOKEN", - Optional: boolPtr(true), - }}, - }, - ) - return job +func (b *Backend) fetchUIPackage(ctx context.Context, namespace string, root string, sourcePath string) ([]byte, error) { + service := serviceName(root, "dev", "webdav") + if b.uiPackageFetcher != nil { + return b.uiPackageFetcher(ctx, namespace, service, sourcePath, b.config.InternalToken) + } + if b.restConfig == nil || b.httpClient == nil { + return nil, fmt.Errorf("kubernetes service proxy is unavailable") + } + segments := strings.Split(sourcePath, "/") + for index := range segments { + segments[index] = url.PathEscape(segments[index]) + } + proxyURL := strings.TrimRight(b.restConfig.Host, "/") + "/api/v1/namespaces/" + namespace + "/services/http:" + service + ":8084/proxy/internal/v1/ui/" + strings.Join(segments, "/") + request, err := http.NewRequestWithContext(ctx, http.MethodGet, proxyURL, nil) + if err != nil { + return nil, err + } + request.Header.Set("Authorization", "Bearer "+b.config.InternalToken) + response, err := b.httpClient.Do(request) + if err != nil { + return nil, err + } + defer response.Body.Close() + if response.StatusCode < http.StatusOK || response.StatusCode >= http.StatusMultipleChoices { + body, _ := io.ReadAll(io.LimitReader(response.Body, 4096)) + return nil, fmt.Errorf("dev service returned %s: %s", response.Status, strings.TrimSpace(string(body))) + } + return io.ReadAll(response.Body) +} + +func (b *Backend) uploadUIPackage(ctx context.Context, content []byte, config uipackage.S3Config) (string, error) { + if b.uiPackageUploader != nil { + return b.uiPackageUploader(ctx, content, config) + } + return uipackage.Upload(ctx, content, config) } diff --git a/internal/uipackage/publish.go b/internal/uipackage/publish.go new file mode 100644 index 00000000..ebe19efb --- /dev/null +++ b/internal/uipackage/publish.go @@ -0,0 +1,60 @@ +package uipackage + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "fmt" + "path" + "strings" + + "github.com/aws/aws-sdk-go-v2/aws" + awscfg "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/credentials" + "github.com/aws/aws-sdk-go-v2/service/s3" +) + +type S3Config struct { + Bucket string + Region string + Endpoint string + KeyPrefix string + AccessKey string + SecretKey string + SessionToken string +} + +// Upload stores content at a content-addressed key and returns its SHA-256. +func Upload(ctx context.Context, content []byte, config S3Config) (string, error) { + if config.Bucket == "" || config.Region == "" { + return "", fmt.Errorf("ui package S3 bucket and region are required") + } + sum := sha256.Sum256(content) + hash := hex.EncodeToString(sum[:]) + cfg, err := awscfg.LoadDefaultConfig(ctx, awscfg.WithRegion(config.Region)) + if err != nil { + return "", err + } + if config.AccessKey != "" || config.SecretKey != "" { + cfg.Credentials = aws.NewCredentialsCache(credentials.NewStaticCredentialsProvider(config.AccessKey, config.SecretKey, config.SessionToken)) + } + client := s3.NewFromConfig(cfg, func(options *s3.Options) { + if config.Endpoint != "" { + options.BaseEndpoint = aws.String(config.Endpoint) + options.UsePathStyle = true + } + }) + key := path.Join(strings.Trim(config.KeyPrefix, "/"), hash, "app.wasm") + _, err = client.PutObject(ctx, &s3.PutObjectInput{ + Bucket: aws.String(config.Bucket), + Key: aws.String(key), + Body: bytes.NewReader(content), + ContentType: aws.String("application/wasm"), + CacheControl: aws.String("public, max-age=31536000, immutable"), + }) + if err != nil { + return "", err + } + return hash, nil +}