From e1235bb894b7d45afc7cca8002df90586f2bee13 Mon Sep 17 00:00:00 2001 From: Jeefos Date: Mon, 25 May 2026 05:30:03 -0400 Subject: [PATCH] file manager --- cmd/kleffd/main.go | 16 +- internal/adapters/in/fileapi/server.go | 475 +++++++++++++++++++++++ internal/adapters/out/platform/client.go | 13 +- internal/app/config/config.go | 6 + 4 files changed, 504 insertions(+), 6 deletions(-) create mode 100644 internal/adapters/in/fileapi/server.go diff --git a/cmd/kleffd/main.go b/cmd/kleffd/main.go index 8dcccf6..2baa51f 100644 --- a/cmd/kleffd/main.go +++ b/cmd/kleffd/main.go @@ -9,6 +9,7 @@ import ( "syscall" "time" + "github.com/kleffio/kleff-daemon/internal/adapters/in/fileapi" "github.com/kleffio/kleff-daemon/internal/adapters/out/db" "github.com/kleffio/kleff-daemon/internal/adapters/out/observability/logging" platformadapter "github.com/kleffio/kleff-daemon/internal/adapters/out/platform" @@ -82,8 +83,16 @@ func main() { } } + // --- File API server --- + fileServer := fileapi.NewServer(cfg.StoragePath, cfg.SharedSecret, cfg.FileAPIPort, daemonLog) + go func() { + if err := fileServer.Start(); err != nil { + daemonLog.Warn("File API server stopped", "error", err) + } + }() + // --- Platform registration + status reporting --- - platformClient := platformadapter.NewClient(cfg.PlatformURL, cfg.SharedSecret, cfg.NodeID, daemonLog) + platformClient := platformadapter.NewClient(cfg.PlatformURL, cfg.SharedSecret, cfg.NodeID, cfg.FileAPIURL, daemonLog) if err := platformClient.RegisterNode(context.Background()); err != nil { daemonLog.Error("Failed to register node with platform", err) os.Exit(1) @@ -114,6 +123,11 @@ func main() { go tailer.Run(ctx) dispatcher.Run(ctx) + + shutCtx, shutCancel := context.WithTimeout(context.Background(), 5*time.Second) + defer shutCancel() + _ = fileServer.Shutdown(shutCtx) + daemonLog.Info("Daemon shutdown complete") } diff --git a/internal/adapters/in/fileapi/server.go b/internal/adapters/in/fileapi/server.go new file mode 100644 index 0000000..99b40f7 --- /dev/null +++ b/internal/adapters/in/fileapi/server.go @@ -0,0 +1,475 @@ +package fileapi + +import ( + "archive/zip" + "context" + "encoding/json" + "fmt" + "io" + "mime" + "net/http" + "os" + "path/filepath" + "sort" + "strings" + "time" + + "github.com/kleffio/kleff-daemon/internal/application/ports" +) + +// Server is a lightweight HTTP file-management API for workload data directories. +// It listens on a configurable port and is authenticated by a shared secret. +type Server struct { + storagePath string + sharedSecret string + logger ports.Logger + srv *http.Server +} + +func NewServer(storagePath, sharedSecret string, port int, logger ports.Logger) *Server { + s := &Server{ + storagePath: storagePath, + sharedSecret: sharedSecret, + logger: logger, + } + + mux := http.NewServeMux() + mux.HandleFunc("GET /v1/{projectID}/{workloadID}/files", s.listOrStat) + mux.HandleFunc("GET /v1/{projectID}/{workloadID}/files/download", s.download) + mux.HandleFunc("GET /v1/{projectID}/{workloadID}/files/export", s.exportZip) + mux.HandleFunc("POST /v1/{projectID}/{workloadID}/files/upload", s.upload) + mux.HandleFunc("POST /v1/{projectID}/{workloadID}/files/rename", s.rename) + mux.HandleFunc("POST /v1/{projectID}/{workloadID}/files/mkdir", s.mkdir) + mux.HandleFunc("POST /v1/{projectID}/{workloadID}/files/import", s.importZip) + mux.HandleFunc("DELETE /v1/{projectID}/{workloadID}/files", s.delete) + + s.srv = &http.Server{ + Addr: fmt.Sprintf(":%d", port), + Handler: s.auth(mux), + ReadTimeout: 30 * time.Second, + WriteTimeout: 60 * time.Second, + } + return s +} + +func (s *Server) Start() error { + s.logger.Info("File API server listening", "addr", s.srv.Addr) + if err := s.srv.ListenAndServe(); err != nil && err != http.ErrServerClosed { + return err + } + return nil +} + +func (s *Server) Shutdown(ctx context.Context) error { + return s.srv.Shutdown(ctx) +} + +// auth middleware validates the Bearer token. +func (s *Server) auth(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if s.sharedSecret == "" { + http.Error(w, "server misconfigured: no shared secret", http.StatusInternalServerError) + return + } + token := strings.TrimPrefix(r.Header.Get("Authorization"), "Bearer ") + if token != s.sharedSecret { + http.Error(w, "unauthorized", http.StatusUnauthorized) + return + } + next.ServeHTTP(w, r) + }) +} + +// workloadRoot returns the absolute root directory for a workload. +// Must match the naming convention used by the docker adapter: kleff_proj_{shortID(projectID)}_{workloadID}. +func (s *Server) workloadRoot(projectID, workloadID string) string { + return filepath.Join(s.storagePath, "kleff_proj_"+shortID(projectID)+"_"+workloadID) +} + +// shortID strips dashes and truncates to 12 chars, matching the docker adapter's convention. +func shortID(s string) string { + s = strings.ReplaceAll(s, "-", "") + if len(s) > 12 { + s = s[:12] + } + return s +} + +// safePath resolves a relative path inside the workload root and rejects traversal. +func safePath(root, rel string) (string, error) { + if rel == "" || rel == "/" { + return root, nil + } + // Clean and make relative + rel = filepath.FromSlash(rel) + rel = filepath.Clean(rel) + if filepath.IsAbs(rel) { + rel = rel[1:] + } + abs := filepath.Join(root, rel) + if !strings.HasPrefix(abs, root+string(filepath.Separator)) && abs != root { + return "", fmt.Errorf("path traversal detected") + } + return abs, nil +} + +type fileEntry struct { + Name string `json:"name"` + Path string `json:"path"` + IsDir bool `json:"is_dir"` + Size int64 `json:"size"` + ModTime time.Time `json:"mod_time"` +} + +func (s *Server) listOrStat(w http.ResponseWriter, r *http.Request) { + root, rel, ok := s.resolve(w, r) + if !ok { + return + } + target, err := safePath(root, rel) + if err != nil { + jsonError(w, err.Error(), http.StatusBadRequest) + return + } + + info, err := os.Stat(target) + if err != nil { + if os.IsNotExist(err) { + // Workload directory doesn't exist yet (server not yet started). + // Return an empty listing rather than 404 so the UI shows an empty file tree. + if target == root { + jsonOK(w, []fileEntry{}) + return + } + jsonError(w, "not found", http.StatusNotFound) + } else { + jsonError(w, err.Error(), http.StatusInternalServerError) + } + return + } + + if !info.IsDir() { + // Return single file stat. + jsonOK(w, fileEntry{ + Name: info.Name(), + Path: rel, + IsDir: false, + Size: info.Size(), + ModTime: info.ModTime().UTC(), + }) + return + } + + entries, err := os.ReadDir(target) + if err != nil { + jsonError(w, err.Error(), http.StatusInternalServerError) + return + } + + sort.Slice(entries, func(i, j int) bool { + if entries[i].IsDir() != entries[j].IsDir() { + return entries[i].IsDir() + } + return entries[i].Name() < entries[j].Name() + }) + + items := make([]fileEntry, 0, len(entries)) + for _, e := range entries { + info, err := e.Info() + if err != nil { + continue + } + entryRel := filepath.ToSlash(filepath.Join(rel, e.Name())) + if !strings.HasPrefix(entryRel, "/") { + entryRel = "/" + entryRel + } + items = append(items, fileEntry{ + Name: e.Name(), + Path: entryRel, + IsDir: e.IsDir(), + Size: info.Size(), + ModTime: info.ModTime().UTC(), + }) + } + jsonOK(w, items) +} + +func (s *Server) download(w http.ResponseWriter, r *http.Request) { + root, rel, ok := s.resolve(w, r) + if !ok { + return + } + target, err := safePath(root, rel) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + info, err := os.Stat(target) + if err != nil { + http.Error(w, "not found", http.StatusNotFound) + return + } + if info.IsDir() { + http.Error(w, "path is a directory", http.StatusBadRequest) + return + } + + ct := mime.TypeByExtension(filepath.Ext(target)) + if ct == "" { + ct = "application/octet-stream" + } + w.Header().Set("Content-Type", ct) + w.Header().Set("Content-Disposition", fmt.Sprintf(`attachment; filename="%s"`, info.Name())) + http.ServeFile(w, r, target) +} + +func (s *Server) upload(w http.ResponseWriter, r *http.Request) { + root, rel, ok := s.resolve(w, r) + if !ok { + return + } + destDir, err := safePath(root, rel) + if err != nil { + jsonError(w, err.Error(), http.StatusBadRequest) + return + } + if err := os.MkdirAll(destDir, 0755); err != nil { + jsonError(w, err.Error(), http.StatusInternalServerError) + return + } + + if err := r.ParseMultipartForm(64 << 20); err != nil { + jsonError(w, "parse multipart: "+err.Error(), http.StatusBadRequest) + return + } + files := r.MultipartForm.File["files"] + if len(files) == 0 { + jsonError(w, "no files in request", http.StatusBadRequest) + return + } + + for _, fh := range files { + src, err := fh.Open() + if err != nil { + jsonError(w, "open upload: "+err.Error(), http.StatusInternalServerError) + return + } + dest := filepath.Join(destDir, filepath.Base(fh.Filename)) + // Prevent traversal via filename. + if !strings.HasPrefix(dest, root) { + src.Close() + jsonError(w, "invalid filename", http.StatusBadRequest) + return + } + f, err := os.Create(dest) + if err != nil { + src.Close() + jsonError(w, "create file: "+err.Error(), http.StatusInternalServerError) + return + } + _, copyErr := io.Copy(f, src) + f.Close() + src.Close() + if copyErr != nil { + jsonError(w, "write file: "+copyErr.Error(), http.StatusInternalServerError) + return + } + } + jsonOK(w, map[string]any{"uploaded": len(files)}) +} + +func (s *Server) rename(w http.ResponseWriter, r *http.Request) { + root := s.workloadRoot(r.PathValue("projectID"), r.PathValue("workloadID")) + var body struct { + From string `json:"from"` + To string `json:"to"` + } + if err := json.NewDecoder(r.Body).Decode(&body); err != nil { + jsonError(w, "decode body: "+err.Error(), http.StatusBadRequest) + return + } + src, err := safePath(root, body.From) + if err != nil { + jsonError(w, "invalid from path", http.StatusBadRequest) + return + } + dst, err := safePath(root, body.To) + if err != nil { + jsonError(w, "invalid to path", http.StatusBadRequest) + return + } + if err := os.Rename(src, dst); err != nil { + jsonError(w, err.Error(), http.StatusInternalServerError) + return + } + jsonOK(w, map[string]any{"ok": true}) +} + +func (s *Server) mkdir(w http.ResponseWriter, r *http.Request) { + root, rel, ok := s.resolve(w, r) + if !ok { + return + } + target, err := safePath(root, rel) + if err != nil { + jsonError(w, err.Error(), http.StatusBadRequest) + return + } + if err := os.MkdirAll(target, 0755); err != nil { + jsonError(w, err.Error(), http.StatusInternalServerError) + return + } + jsonOK(w, map[string]any{"ok": true}) +} + +func (s *Server) delete(w http.ResponseWriter, r *http.Request) { + root, rel, ok := s.resolve(w, r) + if !ok { + return + } + if rel == "" || rel == "/" { + jsonError(w, "cannot delete workload root", http.StatusBadRequest) + return + } + target, err := safePath(root, rel) + if err != nil { + jsonError(w, err.Error(), http.StatusBadRequest) + return + } + if err := os.RemoveAll(target); err != nil { + jsonError(w, err.Error(), http.StatusInternalServerError) + return + } + jsonOK(w, map[string]any{"ok": true}) +} + +func (s *Server) exportZip(w http.ResponseWriter, r *http.Request) { + root, rel, ok := s.resolve(w, r) + if !ok { + return + } + target, err := safePath(root, rel) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + if _, err := os.Stat(target); err != nil { + http.Error(w, "not found", http.StatusNotFound) + return + } + + w.Header().Set("Content-Type", "application/zip") + w.Header().Set("Content-Disposition", `attachment; filename="export.zip"`) + + zw := zip.NewWriter(w) + defer zw.Close() + + _ = filepath.WalkDir(target, func(path string, d os.DirEntry, err error) error { + if err != nil || d.IsDir() { + return nil + } + rel, _ := filepath.Rel(target, path) + fw, err := zw.Create(filepath.ToSlash(rel)) + if err != nil { + return err + } + f, err := os.Open(path) + if err != nil { + return err + } + defer f.Close() + _, err = io.Copy(fw, f) + return err + }) +} + +func (s *Server) importZip(w http.ResponseWriter, r *http.Request) { + root, rel, ok := s.resolve(w, r) + if !ok { + return + } + destDir, err := safePath(root, rel) + if err != nil { + jsonError(w, err.Error(), http.StatusBadRequest) + return + } + if err := os.MkdirAll(destDir, 0755); err != nil { + jsonError(w, err.Error(), http.StatusInternalServerError) + return + } + + if err := r.ParseMultipartForm(256 << 20); err != nil { + jsonError(w, "parse multipart: "+err.Error(), http.StatusBadRequest) + return + } + fhs := r.MultipartForm.File["file"] + if len(fhs) == 0 { + jsonError(w, "no file in request", http.StatusBadRequest) + return + } + fh := fhs[0] + src, err := fh.Open() + if err != nil { + jsonError(w, err.Error(), http.StatusInternalServerError) + return + } + defer src.Close() + + zr, err := zip.NewReader(src.(io.ReaderAt), fh.Size) + if err != nil { + jsonError(w, "invalid zip: "+err.Error(), http.StatusBadRequest) + return + } + + for _, f := range zr.File { + dest := filepath.Join(destDir, filepath.FromSlash(f.Name)) + // Zip-slip prevention. + if !strings.HasPrefix(dest, destDir+string(filepath.Separator)) && dest != destDir { + jsonError(w, "zip slip detected", http.StatusBadRequest) + return + } + if f.FileInfo().IsDir() { + _ = os.MkdirAll(dest, 0755) + continue + } + _ = os.MkdirAll(filepath.Dir(dest), 0755) + rc, err := f.Open() + if err != nil { + jsonError(w, err.Error(), http.StatusInternalServerError) + return + } + out, err := os.Create(dest) + if err != nil { + rc.Close() + jsonError(w, err.Error(), http.StatusInternalServerError) + return + } + _, copyErr := io.Copy(out, rc) + out.Close() + rc.Close() + if copyErr != nil { + jsonError(w, copyErr.Error(), http.StatusInternalServerError) + return + } + } + jsonOK(w, map[string]any{"ok": true}) +} + +// resolve extracts and validates the workload root + "path" query param. +func (s *Server) resolve(w http.ResponseWriter, r *http.Request) (root, rel string, ok bool) { + root = s.workloadRoot(r.PathValue("projectID"), r.PathValue("workloadID")) + rel = r.URL.Query().Get("path") + ok = true + return +} + +func jsonOK(w http.ResponseWriter, v any) { + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(v) +} + +func jsonError(w http.ResponseWriter, msg string, code int) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(code) + _ = json.NewEncoder(w).Encode(map[string]string{"error": msg}) +} diff --git a/internal/adapters/out/platform/client.go b/internal/adapters/out/platform/client.go index 4ca3318..67267d0 100644 --- a/internal/adapters/out/platform/client.go +++ b/internal/adapters/out/platform/client.go @@ -17,16 +17,18 @@ type Client struct { baseURL string bootstrapKey string nodeID string + fileAPIURL string nodeToken string httpClient *http.Client logger ports.Logger } -func NewClient(baseURL, bootstrapKey, nodeID string, logger ports.Logger) *Client { +func NewClient(baseURL, bootstrapKey, nodeID, fileAPIURL string, logger ports.Logger) *Client { return &Client{ baseURL: normalizeBaseURL(baseURL), bootstrapKey: bootstrapKey, nodeID: nodeID, + fileAPIURL: fileAPIURL, httpClient: &http.Client{ Timeout: 10 * time.Second, }, @@ -48,10 +50,11 @@ func (c *Client) RegisterNode(ctx context.Context) error { return fmt.Errorf("platform shared secret is required") } payload := map[string]any{ - "node_id": c.nodeID, - "hostname": c.nodeID, - "region": "local", - "ip_address": "", + "node_id": c.nodeID, + "hostname": c.nodeID, + "region": "local", + "ip_address": "", + "file_api_url": c.fileAPIURL, } body, err := json.Marshal(payload) if err != nil { diff --git a/internal/app/config/config.go b/internal/app/config/config.go index 312131e..5bb733e 100644 --- a/internal/app/config/config.go +++ b/internal/app/config/config.go @@ -35,6 +35,8 @@ type Config struct { PlatformURL string `mapstructure:"platform.url"` SharedSecret string `mapstructure:"shared_secret"` StoragePath string `mapstructure:"storage.path"` + FileAPIPort int `mapstructure:"file_api.port"` + FileAPIURL string `mapstructure:"file_api.url"` } func (c *Config) Validate() error { @@ -83,6 +85,8 @@ func Load() (*Config, error) { v.SetDefault("platform.url", "") v.SetDefault("shared_secret", "") v.SetDefault("storage.path", "/var/lib/kleffd/servers") + v.SetDefault("file_api.port", 8083) + v.SetDefault("file_api.url", "") v.SetEnvPrefix("kleff") v.SetEnvKeyReplacer(strings.NewReplacer(".", "_")) @@ -129,6 +133,8 @@ func Load() (*Config, error) { cfg.PlatformURL = v.GetString("platform.url") cfg.SharedSecret = v.GetString("shared_secret") cfg.StoragePath = v.GetString("storage.path") + cfg.FileAPIPort = v.GetInt("file_api.port") + cfg.FileAPIURL = v.GetString("file_api.url") if err := cfg.Validate(); err != nil { return nil, fmt.Errorf("configuration validation failed: %w", err)