diff --git a/internal/admin/router.go b/internal/admin/router.go
index 6169bf5..dd2e7fc 100644
--- a/internal/admin/router.go
+++ b/internal/admin/router.go
@@ -7,6 +7,8 @@ import (
"errors"
"log/slog"
"net/http"
+ "os"
+ "path/filepath"
"strconv"
"github.com/gin-gonic/gin"
@@ -15,11 +17,12 @@ import (
"spark-mcp-go/internal/cluster"
"spark-mcp-go/internal/middleware"
"spark-mcp-go/internal/storage"
+ "spark-mcp-go/internal/uploads"
)
// Mount attaches the /admin sub-router to r, protecting every route with
// bearer-token admin authentication.
-func Mount(r *gin.Engine, repo *storage.ClusterRepo, auditRepo *audit.Repo, adminTokens []string) {
+func Mount(r *gin.Engine, repo *storage.ClusterRepo, uploadRepo *storage.UploadRepo, auditRepo *audit.Repo, uploadStore *uploads.Store, adminTokens []string) {
// HTML page is public — its own modal prompts for the token.
// All other /admin/* endpoints (API + OpenAPI docs) still require auth.
gPublic := r.Group("/admin")
@@ -33,6 +36,9 @@ func Mount(r *gin.Engine, repo *storage.ClusterRepo, auditRepo *audit.Repo, admi
g.PUT("/clusters/:id", updateCluster(repo, auditRepo))
g.DELETE("/clusters/:id", deleteCluster(repo, auditRepo))
g.GET("/audit", listAudit(auditRepo))
+ g.GET("/uploads", uploadsWebHandler)
+ g.GET("/uploads/api", listUploads(uploadRepo))
+ g.DELETE("/uploads/:id", deleteUpload(uploadRepo, auditRepo, uploadStore))
g.GET("/docs", DocsHandler)
g.GET("/docs/spec", OpenAPISpecHandler)
}
@@ -275,6 +281,66 @@ func deleteCluster(repo *storage.ClusterRepo, auditRepo *audit.Repo) gin.Handler
}
}
+func listUploads(repo *storage.UploadRepo) gin.HandlerFunc {
+ return func(c *gin.Context) {
+ search := c.Query("search")
+ limitStr := c.DefaultQuery("limit", "100")
+ limit, err := strconv.Atoi(limitStr)
+ if err != nil {
+ c.JSON(http.StatusBadRequest, gin.H{"error": "invalid limit"})
+ return
+ }
+
+ uploads, err := repo.List(c.Request.Context(), search, limit)
+ if err != nil {
+ c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
+ return
+ }
+ c.JSON(http.StatusOK, uploads)
+ }
+}
+
+func deleteUpload(repo *storage.UploadRepo, auditRepo *audit.Repo, store *uploads.Store) gin.HandlerFunc {
+ return func(c *gin.Context) {
+ id := c.Param("id")
+ ctx := c.Request.Context()
+
+ meta, err := repo.Get(ctx, id)
+ if err != nil && !errors.Is(err, storage.ErrNotFound) {
+ c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
+ return
+ }
+
+ if err := repo.Delete(ctx, id); err != nil && !errors.Is(err, storage.ErrNotFound) {
+ c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
+ return
+ }
+
+ // Always attempt to remove the on-disk files; the DB is only an index.
+ if store != nil {
+ dataPath := filepath.Join(store.Root, id)
+ _ = os.Remove(dataPath)
+ _ = os.Remove(dataPath + ".meta.json")
+ }
+
+ if meta.FileID != "" {
+ details, _ := audit.MarshalDetails(map[string]any{
+ "file_id": id,
+ "name": meta.Name,
+ "size": meta.Size,
+ "sha256": meta.SHA256,
+ })
+ _ = auditRepo.Insert(ctx, &audit.Entry{
+ Actor: actor(c),
+ Action: audit.ActionUploadDelete,
+ ClusterID: id,
+ Details: details,
+ })
+ }
+
+ c.Status(http.StatusNoContent)
+ }
+}
func listAudit(auditRepo *audit.Repo) gin.HandlerFunc {
return func(c *gin.Context) {
limitStr := c.DefaultQuery("limit", "100")
diff --git a/internal/admin/router_test.go b/internal/admin/router_test.go
index db397f6..7e824d0 100644
--- a/internal/admin/router_test.go
+++ b/internal/admin/router_test.go
@@ -7,6 +7,7 @@ import (
"io"
"net/http"
"net/http/httptest"
+ "path/filepath"
"strings"
"testing"
@@ -15,6 +16,7 @@ import (
"spark-mcp-go/internal/audit"
"spark-mcp-go/internal/cluster"
"spark-mcp-go/internal/storage"
+ "spark-mcp-go/internal/uploads"
)
func init() {
@@ -29,7 +31,11 @@ func newTestAdmin(t *testing.T) (*gin.Engine, *storage.ClusterRepo, *audit.Repo)
}
t.Cleanup(func() { _ = db.Close() })
r := gin.New()
- Mount(r, db.Clusters(), audit.NewRepo(db), []string{"good-token"})
+ uploadStore, err := uploads.New(filepath.Join(t.TempDir(), "uploads"))
+ if err != nil {
+ t.Fatalf("new upload store: %v", err)
+ }
+ Mount(r, db.Clusters(), db.Uploads(), audit.NewRepo(db), &uploadStore, []string{"good-token"})
return r, db.Clusters(), audit.NewRepo(db)
}
diff --git a/internal/admin/uploads_test.go b/internal/admin/uploads_test.go
new file mode 100644
index 0000000..657eb39
--- /dev/null
+++ b/internal/admin/uploads_test.go
@@ -0,0 +1,154 @@
+package admin
+
+import (
+ "context"
+ "encoding/json"
+ "errors"
+ "net/http"
+ "os"
+ "path/filepath"
+ "strings"
+ "testing"
+ "time"
+
+ "github.com/gin-gonic/gin"
+
+ "spark-mcp-go/internal/audit"
+ "spark-mcp-go/internal/storage"
+ "spark-mcp-go/internal/uploads"
+)
+
+func newTestAdminUploads(t *testing.T) (*gin.Engine, *storage.UploadRepo, *audit.Repo, *uploads.Store) {
+ t.Helper()
+ db, err := storage.Open(":memory:")
+ if err != nil {
+ t.Fatalf("open: %v", err)
+ }
+ t.Cleanup(func() { _ = db.Close() })
+
+ uploadStore, err := uploads.New(filepath.Join(t.TempDir(), "uploads"))
+ if err != nil {
+ t.Fatalf("new upload store: %v", err)
+ }
+
+ r := gin.New()
+ Mount(r, db.Clusters(), db.Uploads(), audit.NewRepo(db), &uploadStore, []string{"good-token"})
+ return r, db.Uploads(), audit.NewRepo(db), &uploadStore
+}
+
+func TestUploadsWeb_RequiresAuth(t *testing.T) {
+ r, _, _, _ := newTestAdminUploads(t)
+
+ w := doReq(t, r, "GET", "/admin/uploads", "", nil)
+ if w.Code != http.StatusUnauthorized {
+ t.Errorf("GET /admin/uploads without auth: got %d, want %d", w.Code, http.StatusUnauthorized)
+ }
+
+ w = doReq(t, r, "GET", "/admin/uploads", "good-token", nil)
+ if w.Code != http.StatusOK {
+ t.Errorf("GET /admin/uploads with auth: got %d, want %d", w.Code, http.StatusOK)
+ }
+ ct := w.Header().Get("Content-Type")
+ if !strings.Contains(ct, "text/html") {
+ t.Errorf("Content-Type = %q, want text/html", ct)
+ }
+ body := w.Body.String()
+ if !strings.Contains(body, "href='/admin/uploads'") && !strings.Contains(body, "href=\"/admin/uploads\"") {
+ t.Errorf("body missing uploads nav link")
+ }
+ if !strings.Contains(body, "uploads-body") {
+ t.Errorf("body missing uploads table")
+ }
+}
+
+func TestListUploads(t *testing.T) {
+ r, repo, _, _ := newTestAdminUploads(t)
+ ctx := context.Background()
+
+ uploadedAt := time.Unix(0, time.Now().UnixNano())
+ if err := repo.Create(ctx, "00000000000000000000000000000001", "alpha.txt", 10, "deadbeef", uploadedAt); err != nil {
+ t.Fatalf("create upload: %v", err)
+ }
+ if err := repo.Create(ctx, "00000000000000000000000000000002", "beta.txt", 20, "cafebabe", uploadedAt.Add(time.Second)); err != nil {
+ t.Fatalf("create upload: %v", err)
+ }
+
+ w := doReq(t, r, "GET", "/admin/uploads/api", "good-token", nil)
+ if w.Code != http.StatusOK {
+ t.Fatalf("got status %d, want %d", w.Code, http.StatusOK)
+ }
+ var list []map[string]any
+ if err := json.Unmarshal(w.Body.Bytes(), &list); err != nil {
+ t.Fatalf("unmarshal list: %v", err)
+ }
+ if len(list) != 2 {
+ t.Errorf("got %d uploads, want 2", len(list))
+ }
+
+ w = doReq(t, r, "GET", "/admin/uploads/api?search=alpha", "good-token", nil)
+ if w.Code != http.StatusOK {
+ t.Fatalf("search: got status %d", w.Code)
+ }
+ if err := json.Unmarshal(w.Body.Bytes(), &list); err != nil {
+ t.Fatalf("unmarshal search: %v", err)
+ }
+ if len(list) != 1 {
+ t.Errorf("search result: got %d, want 1", len(list))
+ }
+}
+
+func TestDeleteUpload(t *testing.T) {
+ r, repo, auditRepo, store := newTestAdminUploads(t)
+ ctx := context.Background()
+
+ id := "00000000000000000000000000000003"
+ dataPath := filepath.Join(store.Root, id)
+ metaPath := dataPath + ".meta.json"
+
+ if err := os.WriteFile(dataPath, []byte("payload"), 0o640); err != nil {
+ t.Fatalf("create data file: %v", err)
+ }
+ sc := map[string]any{
+ "name": "delete-me.txt",
+ "size": 7,
+ "sha256": "feedface",
+ "uploaded_at": time.Now().Format(time.RFC3339Nano),
+ }
+ b, _ := json.Marshal(sc)
+ if err := os.WriteFile(metaPath, b, 0o600); err != nil {
+ t.Fatalf("create sidecar: %v", err)
+ }
+ if err := repo.Create(ctx, id, "delete-me.txt", 7, "feedface", time.Now()); err != nil {
+ t.Fatalf("create db row: %v", err)
+ }
+
+ w := doReq(t, r, "DELETE", "/admin/uploads/"+id, "good-token", nil)
+ if w.Code != http.StatusNoContent {
+ t.Fatalf("delete: got status %d, want %d: %s", w.Code, http.StatusNoContent, w.Body.String())
+ }
+
+ if _, err := repo.Get(ctx, id); !errors.Is(err, storage.ErrNotFound) {
+ t.Errorf("db row still exists: %v", err)
+ }
+ if _, err := os.Stat(dataPath); !os.IsNotExist(err) {
+ t.Errorf("data file still exists")
+ }
+ if _, err := os.Stat(metaPath); !os.IsNotExist(err) {
+ t.Errorf("sidecar still exists")
+ }
+
+ entries, err := auditRepo.List(ctx, 100)
+ if err != nil {
+ t.Fatalf("list audit: %v", err)
+ }
+ found := false
+ for _, e := range entries {
+ if e.Action == audit.ActionUploadDelete && e.ClusterID == id {
+ found = true
+ break
+ }
+ }
+ if !found {
+ t.Errorf("missing upload.delete audit entry")
+ }
+}
diff --git a/internal/admin/web.go b/internal/admin/web.go
index 515776f..5d8e02d 100644
--- a/internal/admin/web.go
+++ b/internal/admin/web.go
@@ -7,7 +7,7 @@ import (
"github.com/gin-gonic/gin"
)
-//go:embed web/cluster.html
+//go:embed web
var webFS embed.FS
func webHandler(c *gin.Context) {
@@ -18,3 +18,12 @@ func webHandler(c *gin.Context) {
}
c.Data(http.StatusOK, "text/html; charset=utf-8", data)
}
+
+func uploadsWebHandler(c *gin.Context) {
+ data, err := webFS.ReadFile("web/uploads.html")
+ if err != nil {
+ c.String(http.StatusInternalServerError, "uploads.html: %v", err)
+ return
+ }
+ c.Data(http.StatusOK, "text/html; charset=utf-8", data)
+}
diff --git a/internal/admin/web/cluster.html b/internal/admin/web/cluster.html
index f6b13d5..d0d94d4 100644
--- a/internal/admin/web/cluster.html
+++ b/internal/admin/web/cluster.html
@@ -34,6 +34,10 @@
flex-wrap: wrap;
}
header h1 { margin: 0; font-size: 1.1rem; color: var(--accent); }
+ nav { display: flex; gap: .75rem; }
+ nav a { color: var(--muted); text-decoration: none; }
+ nav a:hover { color: var(--text); }
+ nav a.active { color: var(--accent); }
.auth { display: flex; align-items: center; gap: .5rem; flex-wrap: wrap; }
input[type='text'], input[type='password'], input[type='number'], select, textarea {
background: var(--bg);
@@ -111,7 +115,13 @@
- spark-mcp-go Admin
+
diff --git a/internal/admin/web/uploads.html b/internal/admin/web/uploads.html
new file mode 100644
index 0000000..acd4de2
--- /dev/null
+++ b/internal/admin/web/uploads.html
@@ -0,0 +1,377 @@
+
+
+
+
+
+
Uploads - spark-mcp-go Admin
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+ | File ID |
+ Name |
+ Size (KB) |
+ Uploaded At |
+ SHA256 |
+ Actions |
+
+
+
+ | Loading... |
+
+
+
+
+
+
+
+
+
+
diff --git a/internal/storage/upload_repo.go b/internal/storage/upload_repo.go
index 700effb..e693a9f 100644
--- a/internal/storage/upload_repo.go
+++ b/internal/storage/upload_repo.go
@@ -14,11 +14,11 @@ import (
// The .meta.json sidecar remains the source of truth; this struct is the
// queryable mirror used by the admin UI.
type UploadMeta struct {
- FileID string
- Name string
- Size int64
- SHA256 string
- UploadedAt time.Time
+ FileID string `json:"file_id"`
+ Name string `json:"name"`
+ Size int64 `json:"size"`
+ SHA256 string `json:"sha256"`
+ UploadedAt time.Time `json:"uploaded_at"`
}
// UploadRepo provides CRUD operations for upload file index records.
diff --git a/main.go b/main.go
index c0275d0..f311459 100644
--- a/main.go
+++ b/main.go
@@ -84,6 +84,19 @@ func run() error {
logger.Info("uploads.sweep", "deleted", n)
}
+ backfillCtx, backfillCancel := context.WithTimeout(context.Background(), startupSweepTimeout)
+ n, err = uploadStore.Backfill(backfillCtx, db.Uploads())
+ backfillCancel()
+ if errors.Is(err, context.DeadlineExceeded) {
+ logger.Warn("uploads.backfill.timeout", "inserted", n, "err", err)
+ } else if err != nil {
+ logger.Error("uploads.backfill", "err", err)
+ } else {
+ logger.Info("uploads.backfill", "inserted", n)
+ }
+
+ uploadStore.SetRepo(db.Uploads())
+
gin.SetMode(cfg.GinMode)
r := gin.New()
r.Use(gin.Recovery())
@@ -92,10 +105,11 @@ func run() error {
c.JSON(http.StatusOK, gin.H{"ok": true, "version": "0.0.0"})
})
- admin.Mount(r, db.Clusters(), audit.NewRepo(db), cfg.AdminTokens)
+ admin.Mount(r, db.Clusters(), db.Uploads(), audit.NewRepo(db), &uploadStore, cfg.AdminTokens)
mcpHandler, err := mcpsrv.Handler(&tools.Deps{
Logger: logger.With("component", "mcp"),
+ AuditRepo: audit.NewRepo(db),
ClusterRepo: db.Clusters(),
SparkSubmitTimeout: cfg.SparkSubmitTimeout,
HTTPClient: httpclient.New(httpclient.Config{