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"}}
Pipeline diagnostics Partial artifacts are not final scientific results until a reducer is implemented.
New similarity-search check| ID | Workload | Status | Progress | Created |
|---|---|---|---|---|
| {{.ID}} | {{.Workload}} | {{.Status}} | {{.Completed}} / {{.Total}} | {{time .CreatedAt}} |
| No jobs yet. | ||||
| Name | Status | Capabilities | Last heartbeat |
|---|---|---|---|
| {{.Name}} | {{.Status}} | {{range .Capabilities}}{{.}} {{end}} | {{time .LastHeartbeatAt}} |
| No workers registered. | |||
Pipeline check: downloaded partial_result files are diagnostic, not final reduced outputs.
| Chunk | Status | Attempt | Lease | Error |
|---|---|---|---|---|
| {{.ChunkIndex}} | {{.Status}} | {{.Attempt}} / {{.MaxAttempts}} | {{.LeaseOwner}} {{if .LeaseExpiresAt}}until {{time .LeaseExpiresAt}}{{end}} | {{.ErrorCode}} {{.ErrorMessage}} |
| Kind | File | Size | Checksum | |
|---|---|---|---|---|
| {{.Kind}}{{if .Diagnostic}} (diagnostic){{end}} | {{.Filename}} | {{.SizeBytes}} | {{.SHA256}} | {{if .Downloadable}}Download{{end}} |
This creates a diagnostic shard job. Use query_smiles; global multi-shard reduction is not available yet.