Add local operator web interface
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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).
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
{{define "dashboard.html"}}<!doctype html><html><head><meta charset="utf-8"><title>SciMesh</title><style>body{font:16px system-ui;margin:2rem;max-width:1100px}table{border-collapse:collapse;width:100%;margin:1rem 0}td,th{padding:.5rem;border-bottom:1px solid #ddd;text-align:left}.badge{padding:.2rem .5rem;background:#eef;border-radius:.4rem}a.button{display:inline-block;padding:.6rem .9rem;background:#173b66;color:#fff;border-radius:.3rem;text-decoration:none}</style></head><body><h1>SciMesh operator UI</h1><p><span class="badge">Pipeline diagnostics</span> Partial artifacts are not final scientific results until a reducer is implemented.</p><a class="button" href="/ui/jobs/new">New similarity-search check</a><h2>Recent jobs</h2><table><tr><th>ID</th><th>Workload</th><th>Status</th><th>Progress</th><th>Created</th></tr>{{range .Jobs}}<tr><td><a href="/ui/jobs/{{.ID}}">{{.ID}}</a></td><td>{{.Workload}}</td><td>{{.Status}}</td><td>{{.Completed}} / {{.Total}}</td><td>{{time .CreatedAt}}</td></tr>{{else}}<tr><td colspan="5">No jobs yet.</td></tr>{{end}}</table><h2>Workers</h2><table><tr><th>Name</th><th>Status</th><th>Capabilities</th><th>Last heartbeat</th></tr>{{range .Workers}}<tr><td>{{.Name}}</td><td>{{.Status}}</td><td>{{range .Capabilities}}<code>{{.}}</code> {{end}}</td><td>{{time .LastHeartbeatAt}}</td></tr>{{else}}<tr><td colspan="4">No workers registered.</td></tr>{{end}}</table></body></html>{{end}}
|
||||
@@ -0,0 +1 @@
|
||||
{{define "job.html"}}<!doctype html><html><head><meta charset="utf-8"><title>SciMesh job</title><style>body{font:16px system-ui;margin:2rem;max-width:1100px}table{border-collapse:collapse;width:100%;margin:1rem 0}td,th{padding:.5rem;border-bottom:1px solid #ddd;text-align:left}.warn{background:#fff4d6;padding:.7rem}</style></head><body><p><a href="/ui">← Dashboard</a></p><h1>Job {{.ID}}</h1><p class="warn">Pipeline check: downloaded <code>partial_result</code> files are diagnostic, not final reduced outputs.</p><dl><dt>Workload</dt><dd>{{.Workload}}</dd><dt>Status</dt><dd id="status">{{.Status}}</dd><dt>Progress</dt><dd id="progress">{{.Completed}} / {{.Total}} completed; {{.Failed}} failed</dd></dl><h2>Tasks</h2><table><tr><th>Chunk</th><th>Status</th><th>Attempt</th><th>Lease</th><th>Error</th></tr>{{range .Tasks}}<tr><td>{{.ChunkIndex}}</td><td>{{.Status}}</td><td>{{.Attempt}} / {{.MaxAttempts}}</td><td>{{.LeaseOwner}} {{if .LeaseExpiresAt}}until {{time .LeaseExpiresAt}}{{end}}</td><td>{{.ErrorCode}} {{.ErrorMessage}}</td></tr>{{end}}</table><h2>Coordinator artifacts</h2><table><tr><th>Kind</th><th>File</th><th>Size</th><th>Checksum</th><th></th></tr>{{range .Artifacts}}<tr><td>{{.Kind}}{{if .Diagnostic}} (diagnostic){{end}}</td><td>{{.Filename}}</td><td>{{.SizeBytes}}</td><td><code>{{.SHA256}}</code></td><td>{{if .Downloadable}}<a href="/ui/jobs/{{$.ID}}/artifacts/{{.ID}}">Download</a>{{end}}</td></tr>{{end}}</table><script>const id={{printf "%q" .ID}};setInterval(async()=>{try{const r=await fetch('/ui/api/jobs/'+id);if(!r.ok)return;const d=await r.json();document.querySelector('#status').textContent=d.status;document.querySelector('#progress').textContent=d.completed+' / '+d.total+' completed; '+d.failed+' failed'}catch(_){ }},2000)</script></body></html>{{end}}
|
||||
@@ -0,0 +1 @@
|
||||
{{define "new-job.html"}}<!doctype html><html><head><meta charset="utf-8"><title>New SciMesh run</title><style>body{font:16px system-ui;margin:2rem;max-width:760px}label{display:block;margin:.8rem 0}input{width:100%;padding:.45rem}button{padding:.6rem 1rem}#error{color:#a00}</style></head><body><p><a href="/ui">← Dashboard</a></p><h1>New similarity-search pipeline check</h1><p>This creates a diagnostic shard job. Use <code>query_smiles</code>; global multi-shard reduction is not available yet.</p><form id="run"><label>ChEMBL TSV <input type="file" name="file" required accept=".tsv,.txt,text/tab-separated-values"></label><label>Query SMILES <input name="query_smiles" required maxlength="200" value="CCO"></label><label>Top k <input name="top_k" type="number" min="1" max="100000" value="20" required></label><label>Rows per shard <input name="chunk_rows" type="number" min="1" max="100000" value="1000" required></label><button>Upload and create job</button></form><p id="error" role="alert"></p><script>document.querySelector('#run').addEventListener('submit',async e=>{e.preventDefault();const f=new FormData(e.target);const p={query_smiles:f.get('query_smiles'),top_k:Number(f.get('top_k')),progress_every:0};f.set('workload','similarity-search');f.set('parameters',JSON.stringify(p));try{const r=await fetch('/ui/api/jobs/upload',{method:'POST',body:f});const d=await r.json();if(!r.ok)throw Error(d.error||'Upload failed');location.href='/ui/jobs/'+d.job_id}catch(err){document.querySelector('#error').textContent=err.message}})</script></body></html>{{end}}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
Reference in New Issue
Block a user