From e83e0b5e1f354b8361637a603c515a5b0c609737 Mon Sep 17 00:00:00 2001 From: Emil Date: Thu, 23 Jul 2026 22:08:58 +0300 Subject: [PATCH] Add local operator web interface --- STATUS.md | 7 +- coordinator/.env.example | 4 + coordinator/README.md | 12 ++ coordinator/cmd/coordinator/main.go | 4 +- coordinator/docker-compose.yml | 2 + coordinator/internal/infra/config.go | 6 + coordinator/internal/infra/config_test.go | 34 ++++ coordinator/internal/memstore/ui_read.go | 102 ++++++++++ .../storage/postgres/integration_test.go | 43 +++-- .../internal/storage/postgres/ui_read_repo.go | 123 +++++++++++++ .../internal/transport/http/middleware.go | 42 +++++ coordinator/internal/transport/http/server.go | 19 +- .../internal/transport/http/server_test.go | 120 +++++++++++- .../transport/http/templates/dashboard.html | 1 + .../transport/http/templates/job.html | 1 + .../transport/http/templates/new-job.html | 1 + coordinator/internal/transport/http/ui.go | 121 ++++++++++++ coordinator/internal/usecase/ui.go | 174 ++++++++++++++++++ 18 files changed, 792 insertions(+), 24 deletions(-) create mode 100644 coordinator/internal/infra/config_test.go create mode 100644 coordinator/internal/memstore/ui_read.go create mode 100644 coordinator/internal/storage/postgres/ui_read_repo.go create mode 100644 coordinator/internal/transport/http/templates/dashboard.html create mode 100644 coordinator/internal/transport/http/templates/job.html create mode 100644 coordinator/internal/transport/http/templates/new-job.html create mode 100644 coordinator/internal/transport/http/ui.go create mode 100644 coordinator/internal/usecase/ui.go diff --git a/STATUS.md b/STATUS.md index eb33b5b..f6a7a23 100644 --- a/STATUS.md +++ b/STATUS.md @@ -36,7 +36,7 @@ Docker PostgreSQL stack on 2026-07-23. | CTX-08 Distributed similarity-search | Not started | Local reference exists. | | CTX-09 Reducer and final-result API | Not started | Depends on CTX-07 and CTX-08. | | CTX-10 Distributed similarity-graph | Not started | Local reference exists. | -| CTX-11 Dashboard/operator view | Not started | Deferred until API and reducer work. | +| CTX-11 Dashboard/operator view | In progress | `feat/web-interface` adds a protected local view: job/task/worker status, dataset upload, diagnostic partial-artifact download, and polling. Final-result reduction remains CTX-09. | | CTX-12 Reliability, security, CI | In progress | Unit, race, PostgreSQL integration, and smoke checks exist; CI hardening remains. | ## Next recommended assignment @@ -46,8 +46,9 @@ reduction boundaries before implementing distributed search or graph execution. ## Known constraints -- Planner/reducer semantics are not implemented; use the local `scimesh` CLI - for complete workload results. +- Planner/reducer semantics are not implemented; the operator UI labels + `partial_result` files as diagnostic and cannot present them as final output. + Use the local `scimesh` CLI for complete workload results. - The worker/coordinator flow currently accepts both underscore API workload names and hyphenated CLI names while the contract is consolidated. - A real-stack worker test uses a small `query_smiles` shard. Resolving a diff --git a/coordinator/.env.example b/coordinator/.env.example index d20a617..4c1e596 100644 --- a/coordinator/.env.example +++ b/coordinator/.env.example @@ -6,6 +6,10 @@ DATABASE_URL=postgres://scimesh:scimesh@localhost:5432/scimesh?sslmode=disable # Shared bearer token every worker must present. Leave empty to disable auth (dev only). WORKER_AUTH_TOKEN=change-me +# Optional local operator UI. Use a separate value; never reuse the worker token. +# When empty, /ui is disabled. +UI_AUTH_TOKEN= + # Logging. LOG_LEVEL: debug|info|warn|error. LOG_FILE empty = stdout only; # set a path to also write a size-rotated file (kept across restarts). LOG_LEVEL=info diff --git a/coordinator/README.md b/coordinator/README.md index b79b671..43239cc 100644 --- a/coordinator/README.md +++ b/coordinator/README.md @@ -68,6 +68,18 @@ make logs # follow the coordinator make down # stop (add down-clean to drop the DB volume) ``` +To enable the local operator UI, set a separate credential before starting: + +```sh +UI_AUTH_TOKEN='local-ui-secret' make up +# Open http://localhost:8080/ui and use any username with this value as password. +``` + +The UI is disabled by default and never accepts the worker bearer token. +It shows recent jobs, task/worker state, and the per-job partial artifacts. +Those files are explicitly diagnostic until the CTX-09 reducer creates a final +result; the UI does not present them as final scientific output. + `up` starts three services in order: Postgres waits until `pg_isready` passes, a one-shot `migrate` container applies the schema and exits, and only then does the coordinator start — so it never queries a database that has no tables. diff --git a/coordinator/cmd/coordinator/main.go b/coordinator/cmd/coordinator/main.go index dd11a56..be84e11 100644 --- a/coordinator/cmd/coordinator/main.go +++ b/coordinator/cmd/coordinator/main.go @@ -65,6 +65,7 @@ func run() error { jobRepo = postgres.NewJobRepo(pool) workerRepo = postgres.NewWorkerRepo(pool) artifactRepo = postgres.NewArtifactRepo(pool) + uiReadRepo = postgres.NewUIReadRepo(pool) ) useCases := httptransport.UseCases{ @@ -79,6 +80,7 @@ func run() error { UploadArtifact: usecase.NewUploadArtifact(taskRepo, artifactRepo, blobStore, clk), DownloadArtifact: usecase.NewDownloadArtifact(artifactRepo, blobStore), GetTaskInput: usecase.NewGetTaskInput(taskRepo, artifactRepo, blobStore), + Dashboard: usecase.NewDashboard(uiReadRepo), } // Background reapers are tracked so shutdown can wait for them. Without this @@ -105,7 +107,7 @@ func run() error { // pool.Ping backs /health: readiness means the database answers, not just // that the process is alive. api := httptransport.NewServer(useCases, log, cfg.RequestTimeout, cfg.HeartbeatInterval, cfg.MaxUploadBytes, pool.Ping) - err = infra.RunServer(ctx, log, cfg.Addr, api.Handler(cfg.Token)) + err = infra.RunServer(ctx, log, cfg.Addr, api.Handler(cfg.Token, cfg.UIToken)) // Shutdown order matters, and defers alone cannot express it (they run // LIFO, so the deferred stop() would fire *after* the wait below). diff --git a/coordinator/docker-compose.yml b/coordinator/docker-compose.yml index 901ac55..2a6cd6d 100644 --- a/coordinator/docker-compose.yml +++ b/coordinator/docker-compose.yml @@ -49,6 +49,8 @@ services: # Host is the service name: compose resolves it on the project network. DATABASE_URL: postgres://${POSTGRES_USER:-scimesh}:${POSTGRES_PASSWORD:-scimesh}@postgres:5432/${POSTGRES_DB:-scimesh}?sslmode=disable WORKER_AUTH_TOKEN: ${WORKER_AUTH_TOKEN:-dev-token} + # Empty disables /ui. Set this separately from the worker token. + UI_AUTH_TOKEN: ${UI_AUTH_TOKEN:-} DB_MAX_CONNS: "10" REQUEST_TIMEOUT: "15s" LEASE_DURATION: "2m" diff --git a/coordinator/internal/infra/config.go b/coordinator/internal/infra/config.go index 8703a27..5414866 100644 --- a/coordinator/internal/infra/config.go +++ b/coordinator/internal/infra/config.go @@ -24,6 +24,8 @@ type Config struct { DatabaseURL string // Shared bearer token workers must present. Empty disables auth (dev only). Token string + // Local operator UI credential. Empty disables the embedded UI entirely. + UIToken string // Minimum log level: debug, info, warn, error. LogLevel string @@ -77,6 +79,7 @@ func LoadConfig() (Config, error) { // COORDINATOR_TOKEN is the contract name; WORKER_AUTH_TOKEN is the // former name, still honoured so existing .env files keep working. Token: getEnv("COORDINATOR_TOKEN", os.Getenv("WORKER_AUTH_TOKEN")), + UIToken: os.Getenv("UI_AUTH_TOKEN"), LogLevel: getEnv("LOG_LEVEL", "info"), LogFile: os.Getenv("LOG_FILE"), StorageDir: getEnv("COORDINATOR_STORAGE_DIR", "./data"), @@ -94,6 +97,9 @@ func LoadConfig() (Config, error) { if cfg.DatabaseURL == "" { return Config{}, fmt.Errorf("DATABASE_URL is required") } + if cfg.UIToken != "" && cfg.Token != "" && cfg.UIToken == cfg.Token { + return Config{}, fmt.Errorf("UI_AUTH_TOKEN must differ from the worker auth token") + } var err error if cfg.DBMaxConns, err = getEnvInt32("DB_MAX_CONNS", cfg.DBMaxConns); err != nil { diff --git a/coordinator/internal/infra/config_test.go b/coordinator/internal/infra/config_test.go new file mode 100644 index 0000000..f9c0945 --- /dev/null +++ b/coordinator/internal/infra/config_test.go @@ -0,0 +1,34 @@ +package infra + +import ( + "path/filepath" + "strings" + "testing" +) + +func TestLoadConfigRejectsSharedUIAndWorkerToken(t *testing.T) { + t.Setenv("ENV_FILE", filepath.Join(t.TempDir(), "missing.env")) + t.Setenv("DATABASE_URL", "postgres://test") + t.Setenv("COORDINATOR_TOKEN", "shared-secret") + t.Setenv("UI_AUTH_TOKEN", "shared-secret") + + _, err := LoadConfig() + if err == nil || !strings.Contains(err.Error(), "must differ") { + t.Fatalf("LoadConfig error = %v, want distinct-token error", err) + } +} + +func TestLoadConfigAllowsDistinctUIAndWorkerTokens(t *testing.T) { + t.Setenv("ENV_FILE", filepath.Join(t.TempDir(), "missing.env")) + t.Setenv("DATABASE_URL", "postgres://test") + t.Setenv("COORDINATOR_TOKEN", "worker-secret") + t.Setenv("UI_AUTH_TOKEN", "ui-secret") + + cfg, err := LoadConfig() + if err != nil { + t.Fatalf("LoadConfig: %v", err) + } + if cfg.Token != "worker-secret" || cfg.UIToken != "ui-secret" { + t.Fatalf("unexpected tokens: %+v", cfg) + } +} diff --git a/coordinator/internal/memstore/ui_read.go b/coordinator/internal/memstore/ui_read.go new file mode 100644 index 0000000..ce36e5c --- /dev/null +++ b/coordinator/internal/memstore/ui_read.go @@ -0,0 +1,102 @@ +package memstore + +import ( + "context" + "sort" + + "github.com/google/uuid" + + "github.com/emil28092005/SciMesh/coordinator/internal/domain" + "github.com/emil28092005/SciMesh/coordinator/internal/usecase" +) + +// UIReadRepo is the in-memory read projection used by HTTP/UI tests. +type UIReadRepo struct { + jobs *JobRepo + tasks *TaskRepo + workers *WorkerRepo + artifacts *ArtifactRepo +} + +func NewUIReadRepo(j *JobRepo, t *TaskRepo, w *WorkerRepo, a *ArtifactRepo) *UIReadRepo { + return &UIReadRepo{j, t, w, a} +} + +var _ usecase.UIReadRepository = (*UIReadRepo)(nil) + +func (r *UIReadRepo) GetJob(ctx context.Context, id uuid.UUID) (*domain.Job, error) { + return r.jobs.Get(ctx, id) +} +func (r *UIReadRepo) ListJobs(_ context.Context, limit int) ([]domain.Job, error) { + if limit < 1 || limit > 100 { + return nil, domain.ErrInvalidInput + } + r.jobs.mu.Lock() + defer r.jobs.mu.Unlock() + out := make([]domain.Job, 0, len(r.jobs.jobs)) + for _, job := range r.jobs.jobs { + out = append(out, *job) + } + sort.Slice(out, func(i, j int) bool { + if out[i].CreatedAt.Equal(out[j].CreatedAt) { + return out[i].ID.String() > out[j].ID.String() + } + return out[i].CreatedAt.After(out[j].CreatedAt) + }) + if len(out) > limit { + out = out[:limit] + } + return out, nil +} +func (r *UIReadRepo) ListTasksByJob(_ context.Context, jobID uuid.UUID) ([]domain.Task, error) { + r.tasks.mu.Lock() + defer r.tasks.mu.Unlock() + out := []domain.Task{} + for _, task := range r.tasks.tasks { + if task.JobID == jobID { + out = append(out, *clone(task)) + } + } + sort.Slice(out, func(i, j int) bool { return out[i].ChunkIndex < out[j].ChunkIndex }) + return out, nil +} +func (r *UIReadRepo) ListWorkers(_ context.Context, limit int) ([]domain.Worker, error) { + if limit < 1 || limit > 100 { + return nil, domain.ErrInvalidInput + } + r.workers.mu.Lock() + defer r.workers.mu.Unlock() + out := []domain.Worker{} + for _, worker := range r.workers.workers { + copy := *worker + copy.Capabilities = append([]string(nil), worker.Capabilities...) + out = append(out, copy) + } + sort.Slice(out, func(i, j int) bool { + if out[i].LastHeartbeatAt.Equal(out[j].LastHeartbeatAt) { + return out[i].ID.String() > out[j].ID.String() + } + return out[i].LastHeartbeatAt.After(out[j].LastHeartbeatAt) + }) + if len(out) > limit { + out = out[:limit] + } + return out, nil +} +func (r *UIReadRepo) ListArtifactsByJob(_ context.Context, jobID uuid.UUID) ([]domain.Artifact, error) { + r.artifacts.mu.Lock() + defer r.artifacts.mu.Unlock() + out := []domain.Artifact{} + for _, artifact := range r.artifacts.arts { + if artifact.JobID == jobID { + out = append(out, *artifact) + } + } + sort.Slice(out, func(i, j int) bool { + if out[i].CreatedAt.Equal(out[j].CreatedAt) { + return out[i].ID.String() < out[j].ID.String() + } + return out[i].CreatedAt.Before(out[j].CreatedAt) + }) + return out, nil +} diff --git a/coordinator/internal/storage/postgres/integration_test.go b/coordinator/internal/storage/postgres/integration_test.go index fd461bd..cb0c610 100644 --- a/coordinator/internal/storage/postgres/integration_test.go +++ b/coordinator/internal/storage/postgres/integration_test.go @@ -137,30 +137,37 @@ func TestConcurrentClaimGivesEachTaskToExactlyOneWorker(t *testing.T) { claimed = make(map[uuid.UUID]string) wg sync.WaitGroup ) - // More workers than tasks, so the surplus must come back empty rather than - // steal an already-leased row. + // More workers than tasks. With SKIP LOCKED, a concurrent caller can + // transiently see no eligible row while every remaining row is locked by a + // different claim statement. Poll briefly, as a real worker does, before + // treating the queue as empty. This verifies the actual contract: tasks are + // unique and all eventually become claimable without lock contention. for i := 0; i < tasks*2; i++ { wg.Add(1) go func(n int) { defer wg.Done() - task, err := repo.ClaimNext(context.Background(), usecase.ClaimFilter{ - Owner: fmt.Sprintf("worker-%d", n), - Now: now, - LeaseUntil: now.Add(time.Minute), - }) - if err != nil { - t.Errorf("claim: %v", err) + for attempt := 0; attempt < 20; attempt++ { + task, err := repo.ClaimNext(context.Background(), usecase.ClaimFilter{ + Owner: fmt.Sprintf("worker-%d", n), + Now: now, + LeaseUntil: now.Add(time.Minute), + }) + if err != nil { + t.Errorf("claim: %v", err) + return + } + if task == nil || task.JobID != job.ID { + time.Sleep(time.Millisecond) + continue + } + mu.Lock() + if prev, dup := claimed[task.ID]; dup { + t.Errorf("task %s handed to both %s and worker-%d", task.ID, prev, n) + } + claimed[task.ID] = fmt.Sprintf("worker-%d", n) + mu.Unlock() return } - if task == nil || task.JobID != job.ID { - return // empty queue, or a task from another test's job - } - mu.Lock() - defer mu.Unlock() - if prev, dup := claimed[task.ID]; dup { - t.Errorf("task %s handed to both %s and worker-%d", task.ID, prev, n) - } - claimed[task.ID] = fmt.Sprintf("worker-%d", n) }(i) } wg.Wait() diff --git a/coordinator/internal/storage/postgres/ui_read_repo.go b/coordinator/internal/storage/postgres/ui_read_repo.go new file mode 100644 index 0000000..4a24ed2 --- /dev/null +++ b/coordinator/internal/storage/postgres/ui_read_repo.go @@ -0,0 +1,123 @@ +package postgres + +import ( + "context" + "fmt" + + sq "github.com/Masterminds/squirrel" + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/emil28092005/SciMesh/coordinator/internal/domain" + "github.com/emil28092005/SciMesh/coordinator/internal/usecase" +) + +// UIReadRepo contains bounded, deterministic read queries for the operator UI. +type UIReadRepo struct{ pool *pgxpool.Pool } + +func NewUIReadRepo(pool *pgxpool.Pool) *UIReadRepo { return &UIReadRepo{pool: pool} } + +var _ usecase.UIReadRepository = (*UIReadRepo)(nil) + +func (r *UIReadRepo) GetJob(ctx context.Context, id uuid.UUID) (*domain.Job, error) { + job, err := NewJobRepo(r.pool).Get(ctx, id) + return job, err +} + +func (r *UIReadRepo) ListJobs(ctx context.Context, limit int) ([]domain.Job, error) { + if limit < 1 || limit > 100 { + return nil, domain.ErrInvalidInput + } + sql, args, err := psql.Select(jobColumns...).From("jobs").OrderBy("created_at DESC", "id DESC").Limit(uint64(limit)).ToSql() + if err != nil { + return nil, err + } + rows, err := conn(ctx, r.pool).Query(ctx, sql, args...) + if err != nil { + return nil, fmt.Errorf("list jobs: %w", err) + } + defer rows.Close() + jobs := make([]domain.Job, 0) + for rows.Next() { + var j domain.Job + var status string + var inputURI *string + if err := rows.Scan(&j.ID, &j.Workload, &inputURI, &j.Parameters, &status, &j.CreatedAt, &j.CompletedAt); err != nil { + return nil, err + } + if inputURI != nil { + j.InputURI = *inputURI + } + j.Status = domain.JobStatus(status) + jobs = append(jobs, j) + } + return jobs, rows.Err() +} + +func (r *UIReadRepo) ListTasksByJob(ctx context.Context, jobID uuid.UUID) ([]domain.Task, error) { + sql, args, err := psql.Select(taskColumns...).From("tasks").Where(sq.Eq{"job_id": jobID}).OrderBy("chunk_index ASC").ToSql() + if err != nil { + return nil, err + } + rows, err := conn(ctx, r.pool).Query(ctx, sql, args...) + if err != nil { + return nil, fmt.Errorf("list tasks: %w", err) + } + defer rows.Close() + tasks := make([]domain.Task, 0) + for rows.Next() { + task, err := scanTask(rows) + if err != nil { + return nil, err + } + tasks = append(tasks, *task) + } + return tasks, rows.Err() +} + +func (r *UIReadRepo) ListWorkers(ctx context.Context, limit int) ([]domain.Worker, error) { + if limit < 1 || limit > 100 { + return nil, domain.ErrInvalidInput + } + sql, args, err := psql.Select(workerColumns...).From("workers").OrderBy("last_heartbeat_at DESC", "id DESC").Limit(uint64(limit)).ToSql() + if err != nil { + return nil, err + } + rows, err := conn(ctx, r.pool).Query(ctx, sql, args...) + if err != nil { + return nil, fmt.Errorf("list workers: %w", err) + } + defer rows.Close() + workers := make([]domain.Worker, 0) + for rows.Next() { + worker, err := scanWorker(rows) + if err != nil { + return nil, err + } + workers = append(workers, *worker) + } + return workers, rows.Err() +} + +func (r *UIReadRepo) ListArtifactsByJob(ctx context.Context, jobID uuid.UUID) ([]domain.Artifact, error) { + sql, args, err := psql.Select(artifactColumns...).From("artifacts").Where(sq.Eq{"job_id": jobID}).OrderBy("created_at ASC", "id ASC").ToSql() + if err != nil { + return nil, err + } + rows, err := conn(ctx, r.pool).Query(ctx, sql, args...) + if err != nil { + return nil, fmt.Errorf("list artifacts: %w", err) + } + defer rows.Close() + artifacts := make([]domain.Artifact, 0) + for rows.Next() { + var a domain.Artifact + var kind string + if err := rows.Scan(&a.ID, &a.JobID, &a.TaskID, &a.Attempt, &kind, &a.Filename, &a.StorageKey, &a.ContentType, &a.SizeBytes, &a.SHA256, &a.CreatedAt); err != nil { + return nil, err + } + a.Kind = domain.ArtifactKind(kind) + artifacts = append(artifacts, a) + } + return artifacts, rows.Err() +} diff --git a/coordinator/internal/transport/http/middleware.go b/coordinator/internal/transport/http/middleware.go index feb3f4f..2808313 100644 --- a/coordinator/internal/transport/http/middleware.go +++ b/coordinator/internal/transport/http/middleware.go @@ -64,6 +64,48 @@ func withAuth(token string) func(http.Handler) http.Handler { } } +// withBasicAuth protects the local operator UI with a credential distinct from +// the worker bearer token. The username is intentionally ignored; the password +// is the configured UI token. Basic Auth is suitable only for localhost or a +// TLS-terminating trusted reverse proxy. +func withBasicAuth(token string) func(http.Handler) http.Handler { + return func(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _, password, ok := r.BasicAuth() + if !ok || subtle.ConstantTimeCompare([]byte(password), []byte(token)) != 1 { + w.Header().Set("WWW-Authenticate", `Basic realm="SciMesh UI", charset="UTF-8"`) + writeJSON(w, http.StatusUnauthorized, errorResponse{Error: "unauthorized", RequestID: requestIDFrom(r.Context())}) + return + } + next.ServeHTTP(w, r) + }) + } +} + +// withSameOrigin rejects browser form/fetch writes initiated by another origin. +// A missing Origin is allowed for direct local tools; authenticated UI pages use +// the browser-supplied Origin header on state-changing requests. +func withSameOrigin(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodGet || r.Method == http.MethodHead || r.Method == http.MethodOptions { + next.ServeHTTP(w, r) + return + } + origin := r.Header.Get("Origin") + if origin != "" { + scheme := "http" + if r.TLS != nil { + scheme = "https" + } + if origin != scheme+"://"+r.Host { + writeJSON(w, http.StatusForbidden, errorResponse{Error: "cross-origin request rejected", RequestID: requestIDFrom(r.Context())}) + return + } + } + next.ServeHTTP(w, r) + }) +} + // statusRecorder captures the status code for the access log. type statusRecorder struct { http.ResponseWriter diff --git a/coordinator/internal/transport/http/server.go b/coordinator/internal/transport/http/server.go index cf7c304..1142836 100644 --- a/coordinator/internal/transport/http/server.go +++ b/coordinator/internal/transport/http/server.go @@ -27,6 +27,7 @@ type UseCases struct { UploadArtifact *usecase.UploadArtifact DownloadArtifact *usecase.DownloadArtifact GetTaskInput *usecase.GetTaskInput + Dashboard *usecase.Dashboard } type Server struct { @@ -54,7 +55,7 @@ func NewServer(uc UseCases, log *slog.Logger, requestTimeout, heartbeatInterval // Handler builds the router. Go 1.22's ServeMux matches on method and path // wildcards, so no third-party router is needed. -func (s *Server) Handler(token string) http.Handler { +func (s *Server) Handler(token string, uiToken ...string) http.Handler { protected := http.NewServeMux() protected.HandleFunc("POST /workers/register", s.handleRegister) protected.HandleFunc("POST /jobs", s.handleCreateJob) @@ -70,6 +71,22 @@ func (s *Server) Handler(token string) http.Handler { mux := http.NewServeMux() mux.HandleFunc("GET /health", s.handleHealth) + if len(uiToken) > 0 && uiToken[0] != "" && s.uc.Dashboard != nil { + ui := http.NewServeMux() + ui.HandleFunc("GET /ui", s.handleUIHome) + ui.HandleFunc("GET /ui/jobs/new", s.handleUINewJob) + ui.HandleFunc("GET /ui/jobs/{job_id}", s.handleUIJob) + ui.HandleFunc("GET /ui/api/jobs/{job_id}", s.handleUIJobJSON) + ui.HandleFunc("POST /ui/api/jobs/upload", s.handleUploadDataset) + ui.HandleFunc("GET /ui/jobs/{job_id}/artifacts/{artifact_id}", s.handleUIArtifactDownload) + mux.Handle("/ui", chain(ui, withRequestID, withAccessLog(s.log), withBasicAuth(uiToken[0]), withSameOrigin)) + mux.Handle("/ui/", chain(ui, withRequestID, withAccessLog(s.log), withBasicAuth(uiToken[0]), withSameOrigin)) + } else { + // More specific than the protected catch-all: UI absence is not an auth + // failure and does not disclose that a UI feature is configured elsewhere. + mux.HandleFunc("/ui", http.NotFound) + mux.HandleFunc("/ui/", http.NotFound) + } mux.Handle("/", chain(protected, withRequestID, // outermost: every response gets an ID, withAccessLog(s.log), // including the 401s below diff --git a/coordinator/internal/transport/http/server_test.go b/coordinator/internal/transport/http/server_test.go index 9a0c0f0..7d5c0f6 100644 --- a/coordinator/internal/transport/http/server_test.go +++ b/coordinator/internal/transport/http/server_test.go @@ -20,6 +20,7 @@ import ( ) const token = "secret" +const uiToken = "ui-secret" type env struct { ts *httptest.Server @@ -27,6 +28,10 @@ type env struct { } func newEnv(t *testing.T, ready func(context.Context) error) *env { + return newEnvWithUIToken(t, ready, uiToken) +} + +func newEnvWithUIToken(t *testing.T, ready func(context.Context) error, configuredUIToken string) *env { t.Helper() tasks := memstore.NewTaskRepo() jobs := memstore.NewJobRepo() @@ -49,9 +54,10 @@ func newEnv(t *testing.T, ready func(context.Context) error) *env { UploadArtifact: usecase.NewUploadArtifact(tasks, arts, blobs, clk), DownloadArtifact: usecase.NewDownloadArtifact(arts, blobs), GetTaskInput: usecase.NewGetTaskInput(tasks, arts, blobs), + Dashboard: usecase.NewDashboard(memstore.NewUIReadRepo(jobs, tasks, work, arts)), } srv := coordhttp.NewServer(uc, slog.New(slog.NewTextHandler(io.Discard, nil)), 5*time.Second, 15*time.Second, 1<<30, ready) - ts := httptest.NewServer(srv.Handler(token)) + ts := httptest.NewServer(srv.Handler(token, configuredUIToken)) t.Cleanup(ts.Close) return &env{ts: ts, blobs: blobs} } @@ -97,6 +103,118 @@ func TestHealthOK(t *testing.T) { } } +func TestUIRequiresDistinctCredentialAndRendersDashboard(t *testing.T) { + e := newEnv(t, healthy) + request := func() *http.Request { + req, _ := http.NewRequestWithContext(context.Background(), "GET", e.ts.URL+"/ui", nil) + return req + } + resp, err := http.DefaultClient.Do(request()) + if err != nil { + t.Fatal(err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusUnauthorized { + t.Fatalf("no UI auth: %d", resp.StatusCode) + } + + req := request() + req.Header.Set("Authorization", "Bearer "+token) + resp, err = http.DefaultClient.Do(req) + if err != nil { + t.Fatal(err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusUnauthorized { + t.Fatalf("worker token authorized UI: %d", resp.StatusCode) + } + + req = request() + req.SetBasicAuth("operator", uiToken) + resp, err = http.DefaultClient.Do(req) + if err != nil { + t.Fatal(err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + t.Fatalf("UI status: %d", resp.StatusCode) + } + body, _ := io.ReadAll(resp.Body) + if !strings.Contains(string(body), "SciMesh operator UI") { + t.Errorf("dashboard body missing title") + } +} + +func TestUIDisabledReturnsNotFound(t *testing.T) { + e := newEnvWithUIToken(t, healthy, "") + resp := e.get(t, "/ui") + defer resp.Body.Close() + if resp.StatusCode != http.StatusNotFound { + t.Fatalf("disabled UI = %d, want 404", resp.StatusCode) + } +} + +func TestUIRejectsCrossOriginUpload(t *testing.T) { + e := newEnv(t, healthy) + req, _ := http.NewRequestWithContext(context.Background(), "POST", e.ts.URL+"/ui/api/jobs/upload", strings.NewReader("dataset=x")) + req.Header.Set("Content-Type", "application/x-www-form-urlencoded") + req.Header.Set("Origin", "https://attacker.example") + req.SetBasicAuth("operator", uiToken) + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatal(err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusForbidden { + t.Fatalf("cross-origin upload = %d, want 403", resp.StatusCode) + } +} + +func TestUIJobAndArtifactAreScopedToTheirJob(t *testing.T) { + e := newEnv(t, healthy) + code, job := e.do(t, "POST", "/jobs", `{"workload":"w","input_uri":"s3://in","chunks":[{"chunk_index":0,"input_uri":"s3://c","input_sha256":"sha"}]}`) + if code != http.StatusCreated { + t.Fatalf("create: %d", code) + } + req, _ := http.NewRequestWithContext(context.Background(), "GET", e.ts.URL+"/ui/jobs/"+job["id"].(string), nil) + req.SetBasicAuth("operator", uiToken) + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatal(err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + t.Fatalf("detail: %d", resp.StatusCode) + } + if got := resp.Header.Get("Content-Security-Policy"); got == "" { + t.Error("missing UI CSP") + } +} + +func TestUIArtifactDownloadRejectsAnotherJobsArtifact(t *testing.T) { + e := newEnv(t, healthy) + code, _ := e.do(t, "POST", "/jobs", `{"workload":"w","input_uri":"s3://in","chunks":[{"chunk_index":0,"input_uri":"s3://c","input_sha256":"sha"}]}`) + if code != http.StatusCreated { + t.Fatalf("first job: %d", code) + } + _, claim := e.do(t, "POST", "/tasks/claim", `{"worker_id":"w1","capabilities":["w"]}`) + artifactID := e.putArtifact(t, claim["task_id"].(string), "w1", int(claim["attempt"].(float64)), "result") + code, second := e.do(t, "POST", "/jobs", `{"workload":"w","input_uri":"s3://in","chunks":[{"chunk_index":0,"input_uri":"s3://c","input_sha256":"sha"}]}`) + if code != http.StatusCreated { + t.Fatalf("second job: %d", code) + } + req, _ := http.NewRequestWithContext(context.Background(), "GET", e.ts.URL+"/ui/jobs/"+second["id"].(string)+"/artifacts/"+artifactID, nil) + req.SetBasicAuth("operator", uiToken) + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatal(err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusNotFound { + t.Errorf("cross-job artifact = %d, want 404", resp.StatusCode) + } +} + func TestHealthUnavailableWhenDBDown(t *testing.T) { e := newEnv(t, func(context.Context) error { return context.DeadlineExceeded }) resp := e.get(t, "/health") diff --git a/coordinator/internal/transport/http/templates/dashboard.html b/coordinator/internal/transport/http/templates/dashboard.html new file mode 100644 index 0000000..2f620b6 --- /dev/null +++ b/coordinator/internal/transport/http/templates/dashboard.html @@ -0,0 +1 @@ +{{define "dashboard.html"}}SciMesh

SciMesh operator UI

Pipeline diagnostics Partial artifacts are not final scientific results until a reducer is implemented.

New similarity-search check

Recent jobs

{{range .Jobs}}{{else}}{{end}}
IDWorkloadStatusProgressCreated
{{.ID}}{{.Workload}}{{.Status}}{{.Completed}} / {{.Total}}{{time .CreatedAt}}
No jobs yet.

Workers

{{range .Workers}}{{else}}{{end}}
NameStatusCapabilitiesLast heartbeat
{{.Name}}{{.Status}}{{range .Capabilities}}{{.}} {{end}}{{time .LastHeartbeatAt}}
No workers registered.
{{end}} diff --git a/coordinator/internal/transport/http/templates/job.html b/coordinator/internal/transport/http/templates/job.html new file mode 100644 index 0000000..44fec9a --- /dev/null +++ b/coordinator/internal/transport/http/templates/job.html @@ -0,0 +1 @@ +{{define "job.html"}}SciMesh job

← Dashboard

Job {{.ID}}

Pipeline check: downloaded partial_result files are diagnostic, not final reduced outputs.

Workload
{{.Workload}}
Status
{{.Status}}
Progress
{{.Completed}} / {{.Total}} completed; {{.Failed}} failed

Tasks

{{range .Tasks}}{{end}}
ChunkStatusAttemptLeaseError
{{.ChunkIndex}}{{.Status}}{{.Attempt}} / {{.MaxAttempts}}{{.LeaseOwner}} {{if .LeaseExpiresAt}}until {{time .LeaseExpiresAt}}{{end}}{{.ErrorCode}} {{.ErrorMessage}}

Coordinator artifacts

{{range .Artifacts}}{{end}}
KindFileSizeChecksum
{{.Kind}}{{if .Diagnostic}} (diagnostic){{end}}{{.Filename}}{{.SizeBytes}}{{.SHA256}}{{if .Downloadable}}Download{{end}}
{{end}} diff --git a/coordinator/internal/transport/http/templates/new-job.html b/coordinator/internal/transport/http/templates/new-job.html new file mode 100644 index 0000000..f01a413 --- /dev/null +++ b/coordinator/internal/transport/http/templates/new-job.html @@ -0,0 +1 @@ +{{define "new-job.html"}}New SciMesh run

← Dashboard

New similarity-search pipeline check

This creates a diagnostic shard job. Use query_smiles; global multi-shard reduction is not available yet.

{{end}} diff --git a/coordinator/internal/transport/http/ui.go b/coordinator/internal/transport/http/ui.go new file mode 100644 index 0000000..dc44edb --- /dev/null +++ b/coordinator/internal/transport/http/ui.go @@ -0,0 +1,121 @@ +package http + +import ( + "embed" + "html/template" + "io" + "mime" + "net/http" + "strconv" + "time" + + "github.com/google/uuid" + + "github.com/emil28092005/SciMesh/coordinator/internal/domain" +) + +//go:embed templates/*.html +var uiFiles embed.FS + +var uiTemplates = template.Must(template.New("ui").Funcs(template.FuncMap{ + "time": func(t time.Time) string { return t.UTC().Format(time.RFC3339) }, +}).ParseFS(uiFiles, "templates/*.html")) + +func (s *Server) renderUI(w http.ResponseWriter, name string, data any) { + w.Header().Set("Content-Type", "text/html; charset=utf-8") + w.Header().Set("Cache-Control", "no-store") + w.Header().Set("X-Content-Type-Options", "nosniff") + w.Header().Set("Content-Security-Policy", "default-src 'self'; style-src 'self' 'unsafe-inline'; script-src 'self' 'unsafe-inline'; base-uri 'none'; frame-ancestors 'none'") + if err := uiTemplates.ExecuteTemplate(w, name, data); err != nil { + s.log.Error("render UI", "err", err) + http.Error(w, "internal error", http.StatusInternalServerError) + } +} + +func (s *Server) handleUIHome(w http.ResponseWriter, r *http.Request) { + ctx, cancel := s.reqCtx(r) + defer cancel() + view, err := s.uc.Dashboard.Overview(ctx, 20) + if err != nil { + s.writeError(w, r, err) + return + } + s.renderUI(w, "dashboard.html", view) +} + +func (s *Server) handleUINewJob(w http.ResponseWriter, r *http.Request) { + s.renderUI(w, "new-job.html", nil) +} + +func (s *Server) uiJobID(w http.ResponseWriter, r *http.Request) (uuid.UUID, bool) { + return s.pathUUID(w, r, "job_id") +} + +func (s *Server) handleUIJob(w http.ResponseWriter, r *http.Request) { + jobID, ok := s.uiJobID(w, r) + if !ok { + return + } + ctx, cancel := s.reqCtx(r) + defer cancel() + view, err := s.uc.Dashboard.JobDetail(ctx, jobID) + if err != nil { + s.writeError(w, r, err) + return + } + s.renderUI(w, "job.html", view) +} + +func (s *Server) handleUIJobJSON(w http.ResponseWriter, r *http.Request) { + jobID, ok := s.uiJobID(w, r) + if !ok { + return + } + ctx, cancel := s.reqCtx(r) + defer cancel() + view, err := s.uc.Dashboard.JobDetail(ctx, jobID) + if err != nil { + s.writeError(w, r, err) + return + } + writeJSON(w, http.StatusOK, view) +} + +func (s *Server) handleUIArtifactDownload(w http.ResponseWriter, r *http.Request) { + jobID, ok := s.uiJobID(w, r) + if !ok { + return + } + artifactID, err := uuid.Parse(r.PathValue("artifact_id")) + if err != nil { + s.writeError(w, r, domain.ErrInvalidInput) + return + } + ctx, cancel := s.reqCtx(r) + defer cancel() + belongs, err := s.uc.Dashboard.ArtifactBelongsToJob(ctx, jobID, artifactID) + if err != nil { + s.writeError(w, r, err) + return + } + if !belongs { + s.writeError(w, r, domain.ErrArtifactNotFound) + return + } + // Reuse the coordinator-owned blob stream after the job-scoped check above. + art, body, err := s.uc.DownloadArtifact.Execute(ctx, artifactID) + if err != nil { + s.writeError(w, r, err) + return + } + defer func() { + if err := body.Close(); err != nil { + s.log.Warn("close downloaded UI artifact", "artifact_id", artifactID, "err", err) + } + }() + w.Header().Set("Content-Type", art.ContentType) + w.Header().Set("Content-Disposition", mime.FormatMediaType("attachment", map[string]string{"filename": art.Filename})) + w.Header().Set("Content-Length", strconv.FormatInt(art.SizeBytes, 10)) + w.Header().Set("X-Checksum-SHA256", art.SHA256) + _, _ = io.Copy(w, body) +} diff --git a/coordinator/internal/usecase/ui.go b/coordinator/internal/usecase/ui.go new file mode 100644 index 0000000..f250480 --- /dev/null +++ b/coordinator/internal/usecase/ui.go @@ -0,0 +1,174 @@ +package usecase + +import ( + "context" + "time" + + "github.com/google/uuid" + + "github.com/emil28092005/SciMesh/coordinator/internal/domain" +) + +// UIReadRepository is a read-only projection source for the local operator UI. +// It intentionally exposes no storage paths or credentials. +type UIReadRepository interface { + GetJob(ctx context.Context, jobID uuid.UUID) (*domain.Job, error) + ListJobs(ctx context.Context, limit int) ([]domain.Job, error) + ListTasksByJob(ctx context.Context, jobID uuid.UUID) ([]domain.Task, error) + ListWorkers(ctx context.Context, limit int) ([]domain.Worker, error) + ListArtifactsByJob(ctx context.Context, jobID uuid.UUID) ([]domain.Artifact, error) +} + +type JobCard struct { + ID string `json:"id"` + Workload string `json:"workload"` + Status string `json:"status"` + CreatedAt time.Time `json:"created_at"` + Total int `json:"total"` + Pending int `json:"pending"` + Leased int `json:"leased"` + Running int `json:"running"` + Completed int `json:"completed"` + Failed int `json:"failed"` +} + +type TaskCard struct { + ID string `json:"id"` + ChunkIndex int `json:"chunk_index"` + Status string `json:"status"` + Attempt int `json:"attempt"` + MaxAttempts int `json:"max_attempts"` + LeaseOwner string `json:"lease_owner,omitempty"` + LeaseExpiresAt *time.Time `json:"lease_expires_at,omitempty"` + ErrorCode string `json:"error_code,omitempty"` + ErrorMessage string `json:"error_message,omitempty"` +} + +type ArtifactCard struct { + ID string `json:"id"` + Kind string `json:"kind"` + Filename string `json:"filename"` + SizeBytes int64 `json:"size_bytes"` + SHA256 string `json:"sha256"` + Downloadable bool `json:"downloadable"` + Diagnostic bool `json:"diagnostic"` +} + +type WorkerCard struct { + ID string `json:"id"` + Name string `json:"name"` + Status string `json:"status"` + Capabilities []string `json:"capabilities"` + LastHeartbeatAt time.Time `json:"last_heartbeat_at"` +} + +type DashboardView struct { + Jobs []JobCard + Workers []WorkerCard +} +type JobDetailView struct { + JobCard + Tasks []TaskCard `json:"tasks"` + Artifacts []ArtifactCard `json:"artifacts"` + FinalResultAvailable bool `json:"final_result_available"` +} + +type Dashboard struct{ read UIReadRepository } + +func NewDashboard(read UIReadRepository) *Dashboard { return &Dashboard{read: read} } + +func (d *Dashboard) Overview(ctx context.Context, limit int) (DashboardView, error) { + jobs, err := d.read.ListJobs(ctx, limit) + if err != nil { + return DashboardView{}, err + } + workers, err := d.read.ListWorkers(ctx, limit) + if err != nil { + return DashboardView{}, err + } + out := DashboardView{Jobs: make([]JobCard, 0, len(jobs)), Workers: make([]WorkerCard, 0, len(workers))} + for _, job := range jobs { + tasks, err := d.read.ListTasksByJob(ctx, job.ID) + if err != nil { + return DashboardView{}, err + } + out.Jobs = append(out.Jobs, jobCard(job, tasks)) + } + for _, worker := range workers { + out.Workers = append(out.Workers, WorkerCard{ID: worker.ID.String(), Name: worker.Name, Status: string(worker.Status), Capabilities: worker.Capabilities, LastHeartbeatAt: worker.LastHeartbeatAt}) + } + return out, nil +} + +func (d *Dashboard) JobDetail(ctx context.Context, jobID uuid.UUID) (JobDetailView, error) { + job, err := d.read.GetJob(ctx, jobID) + if err != nil { + return JobDetailView{}, err + } + tasks, err := d.read.ListTasksByJob(ctx, jobID) + if err != nil { + return JobDetailView{}, err + } + artifacts, err := d.read.ListArtifactsByJob(ctx, jobID) + if err != nil { + return JobDetailView{}, err + } + out := JobDetailView{JobCard: jobCard(*job, tasks), Tasks: make([]TaskCard, 0, len(tasks)), Artifacts: make([]ArtifactCard, 0, len(artifacts))} + for _, task := range tasks { + card := TaskCard{ID: task.ID.String(), ChunkIndex: task.ChunkIndex, Status: string(task.Status), Attempt: task.Attempt, MaxAttempts: task.MaxAttempts, LeaseExpiresAt: task.LeaseExpiresAt} + if task.LeaseOwner != nil { + card.LeaseOwner = *task.LeaseOwner + } + if task.ErrorCode != nil { + card.ErrorCode = *task.ErrorCode + } + if task.ErrorMessage != nil { + card.ErrorMessage = *task.ErrorMessage + } + out.Tasks = append(out.Tasks, card) + } + for _, artifact := range artifacts { + diagnostic := artifact.Kind == domain.ArtifactPartialResult + downloadable := diagnostic || (artifact.Kind == domain.ArtifactFinalResult && out.Status == string(domain.JobCompleted)) + out.Artifacts = append(out.Artifacts, ArtifactCard{ID: artifact.ID.String(), Kind: string(artifact.Kind), Filename: artifact.Filename, SizeBytes: artifact.SizeBytes, SHA256: artifact.SHA256, Downloadable: downloadable, Diagnostic: diagnostic}) + if artifact.Kind == domain.ArtifactFinalResult && downloadable { + out.FinalResultAvailable = true + } + } + return out, nil +} + +func (d *Dashboard) ArtifactBelongsToJob(ctx context.Context, jobID, artifactID uuid.UUID) (bool, error) { + artifacts, err := d.read.ListArtifactsByJob(ctx, jobID) + if err != nil { + return false, err + } + for _, a := range artifacts { + if a.ID == artifactID { + return true, nil + } + } + return false, nil +} + +func jobCard(job domain.Job, tasks []domain.Task) JobCard { + c := JobCard{ID: job.ID.String(), Workload: job.Workload, CreatedAt: job.CreatedAt} + for _, task := range tasks { + c.Total++ + switch task.Status { + case domain.TaskPending: + c.Pending++ + case domain.TaskLeased: + c.Leased++ + case domain.TaskRunning: + c.Running++ + case domain.TaskCompleted: + c.Completed++ + case domain.TaskFailed: + c.Failed++ + } + } + p := domain.JobProgress{Job: job, Total: c.Total, Pending: c.Pending, Leased: c.Leased + c.Running, Done: c.Completed, Failed: c.Failed} + c.Status = string(p.DeriveStatus()) + return c +}