Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 8 additions & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -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..."
Expand Down
29 changes: 24 additions & 5 deletions apps/druid/adapters/cli/client/dev.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package client

import (
"context"
"crypto/subtle"
"encoding/json"
"fmt"
"mime"
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Expand All @@ -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),
Expand Down Expand Up @@ -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()
}
Expand Down Expand Up @@ -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)
}
Expand Down
30 changes: 30 additions & 0 deletions apps/druid/adapters/cli/client/dev_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
13 changes: 10 additions & 3 deletions apps/druid/adapters/cli/daemon.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)")
}

Expand All @@ -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()
Expand Down
33 changes: 3 additions & 30 deletions apps/druid/adapters/cli/ui.go
Original file line number Diff line number Diff line change
@@ -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"
)

Expand Down Expand Up @@ -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) {
Expand Down
3 changes: 2 additions & 1 deletion internal/runtime/docker/ui.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down Expand Up @@ -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
}
Expand Down
26 changes: 18 additions & 8 deletions internal/runtime/kubernetes/backend.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package kubernetes
import (
"context"
"fmt"
"net/http"
"sync"
"time"

Expand All @@ -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 {
Expand Down Expand Up @@ -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),
Expand Down
22 changes: 17 additions & 5 deletions internal/runtime/kubernetes/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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
}
Expand Down
2 changes: 2 additions & 0 deletions internal/runtime/kubernetes/resources.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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}},
}},
Expand Down
Loading
Loading