diff --git a/coordinator/cmd/coordinator/main.go b/coordinator/cmd/coordinator/main.go
index 965fece..90d352f 100644
--- a/coordinator/cmd/coordinator/main.go
+++ b/coordinator/cmd/coordinator/main.go
@@ -77,6 +77,7 @@ type storageDeps struct {
artifactRepo usecase.ArtifactRepository
uiReadRepo usecase.UIReadRepository
adminReadRepo usecase.AdminReadRepository
+ settingsRepo usecase.WorkloadSettingsRepository
taskResultRepo usecase.TaskResultRepository
statsRepo interface {
Counts(ctx context.Context) (tasks, jobs, workers map[string]int, err error)
@@ -159,7 +160,7 @@ func runWithConfig(cfg infra.Config) error {
useCases := httptransport.UseCases{
RegisterWorker: usecase.NewRegisterWorker(workerRepo, clk),
CreateJob: usecase.NewCreateJob(jobRepo, taskRepo, tx, clk),
- SubmitDataset: usecase.NewSubmitDataset(blobStore, artifactRepo, jobRepo, taskRepo, tx, clk, cfg.DefaultMaxAttempts, catalog),
+ SubmitDataset: usecase.NewSubmitDataset(blobStore, artifactRepo, jobRepo, taskRepo, tx, clk, cfg.DefaultMaxAttempts, catalog, deps.settingsRepo),
ClaimTask: usecase.NewClaimTask(taskRepo, jobRepo, workerRepo, tx, clk, cfg.LeaseDuration, catalog),
RenewLease: usecase.NewRenewLease(taskRepo, workerRepo, tx, clk, cfg.LeaseDuration),
CompleteTask: usecase.NewCompleteTask(taskRepo, jobRepo, artifactRepo, workerRepo, taskResultRepo, tx, clk, cfg.QuorumSize, catalog),
@@ -173,16 +174,21 @@ func runWithConfig(cfg infra.Config) error {
GetTaskInput: usecase.NewGetTaskInput(taskRepo, artifactRepo, blobStore),
Dashboard: usecase.NewDashboard(uiReadRepo, catalog),
PreviewArtifact: usecase.NewPreviewArtifact(uiReadRepo, blobStore),
- Admin: usecase.NewAdmin(deps.adminReadRepo, uiReadRepo, usecase.AdminNodeInfo{
- Version: version,
- StartedAt: clk.Now(),
- Binary: executablePath(),
- Addr: cfg.Addr,
- DataDir: cfg.StorageDir,
- DBEngine: cfg.DatabaseEngine,
- PublicURL: cfg.PublicCoordinatorURL,
- Userservice: cfg.UserserviceURL,
- }, deps.ready, clk.Now),
+ Admin: usecase.NewAdmin(deps.adminReadRepo, uiReadRepo, workerRepo, deps.settingsRepo, catalog,
+ usecase.AdminNodeInfo{
+ Version: version,
+ StartedAt: clk.Now(),
+ Binary: executablePath(),
+ Addr: cfg.Addr,
+ DataDir: cfg.StorageDir,
+ DBEngine: cfg.DatabaseEngine,
+ PublicURL: cfg.PublicCoordinatorURL,
+ Userservice: cfg.UserserviceURL,
+ WorkerToken: func() string { return cfg.Token },
+ }, deps.ready, clk.Now).
+ WithAuditLog(log, func(ctx context.Context, action, detail string) {
+ log.Info("admin audit", "action", action, "detail", detail)
+ }),
}
// Background reapers are tracked so shutdown can wait for them. Without this
@@ -264,6 +270,7 @@ func openSQLite(ctx context.Context, cfg infra.Config, log *slog.Logger) (*stora
artifactRepo: sqlite.NewArtifactRepo(db),
uiReadRepo: sqlite.NewUIReadRepo(db),
adminReadRepo: sqlite.NewAdminReadRepo(db),
+ settingsRepo: sqlite.NewWorkloadSettingsRepo(db),
taskResultRepo: sqlite.NewTaskResultRepo(db),
statsRepo: sqlite.NewStatsRepo(db),
ready: func(ctx context.Context) error { return db.PingContext(ctx) },
@@ -287,6 +294,7 @@ func openPostgres(ctx context.Context, cfg infra.Config, log *slog.Logger) (*sto
artifactRepo: postgres.NewArtifactRepo(pool),
uiReadRepo: postgres.NewUIReadRepo(pool),
adminReadRepo: postgres.NewAdminReadRepo(pool),
+ settingsRepo: postgres.NewWorkloadSettingsRepo(pool),
taskResultRepo: postgres.NewTaskResultRepo(pool),
statsRepo: postgres.NewStatsRepo(pool),
ready: func(ctx context.Context) error { return pool.Ping(ctx) },
diff --git a/coordinator/internal/domain/errors.go b/coordinator/internal/domain/errors.go
index fca2979..7e8a825 100644
--- a/coordinator/internal/domain/errors.go
+++ b/coordinator/internal/domain/errors.go
@@ -18,4 +18,5 @@ var (
ErrResultConflict = errors.New("different result already recorded")
ErrInvalidInput = errors.New("invalid input")
ErrTaskNotLeased = errors.New("task is not currently leased")
+ ErrWorkloadDisabled = errors.New("workload is disabled")
)
diff --git a/coordinator/internal/memstore/memstore.go b/coordinator/internal/memstore/memstore.go
index 9c911e4..ace397f 100644
--- a/coordinator/internal/memstore/memstore.go
+++ b/coordinator/internal/memstore/memstore.go
@@ -314,6 +314,17 @@ func (r *WorkerRepo) MarkStaleOffline(ctx context.Context, cutoff time.Time) (in
return n, nil
}
+func (r *WorkerRepo) SetTrust(ctx context.Context, id uuid.UUID, trust domain.WorkerTrust) error {
+ r.mu.Lock()
+ defer r.mu.Unlock()
+ w, ok := r.workers[id]
+ if !ok {
+ return domain.ErrWorkerNotFound
+ }
+ w.TrustLevel = trust
+ return nil
+}
+
// --- ArtifactRepo --------------------------------------------------------
type ArtifactRepo struct {
diff --git a/coordinator/internal/storage/postgres/migrate_test.go b/coordinator/internal/storage/postgres/migrate_test.go
index 7b241f2..f7f3789 100644
--- a/coordinator/internal/storage/postgres/migrate_test.go
+++ b/coordinator/internal/storage/postgres/migrate_test.go
@@ -54,6 +54,8 @@ func expectedMigrationName(version int) string {
return "0012_worker_trust.up.sql"
case 13:
return "0013_task_results.up.sql"
+ case 14:
+ return "0014_workload_settings.up.sql"
default:
return ""
}
diff --git a/coordinator/internal/storage/postgres/migrations/0014_workload_settings.down.sql b/coordinator/internal/storage/postgres/migrations/0014_workload_settings.down.sql
new file mode 100644
index 0000000..b79ba91
--- /dev/null
+++ b/coordinator/internal/storage/postgres/migrations/0014_workload_settings.down.sql
@@ -0,0 +1,5 @@
+BEGIN;
+
+DROP TABLE IF EXISTS workload_settings;
+
+COMMIT;
diff --git a/coordinator/internal/storage/postgres/migrations/0014_workload_settings.up.sql b/coordinator/internal/storage/postgres/migrations/0014_workload_settings.up.sql
new file mode 100644
index 0000000..2975b96
--- /dev/null
+++ b/coordinator/internal/storage/postgres/migrations/0014_workload_settings.up.sql
@@ -0,0 +1,11 @@
+BEGIN;
+
+-- Per-workload enable/disable. Absence of a row means "enabled" (the catalog
+-- default); a row only exists once an admin flipped a workload off or back on.
+CREATE TABLE workload_settings (
+ workload text NOT NULL PRIMARY KEY,
+ enabled boolean NOT NULL,
+ updated_at timestamptz NOT NULL DEFAULT now()
+);
+
+COMMIT;
diff --git a/coordinator/internal/storage/postgres/worker_repo.go b/coordinator/internal/storage/postgres/worker_repo.go
index 5f75183..e3d2bc5 100644
--- a/coordinator/internal/storage/postgres/worker_repo.go
+++ b/coordinator/internal/storage/postgres/worker_repo.go
@@ -91,6 +91,24 @@ func (r *WorkerRepo) MarkStaleOffline(ctx context.Context, cutoff time.Time) (in
return tag.RowsAffected(), nil
}
+func (r *WorkerRepo) SetTrust(ctx context.Context, id uuid.UUID, trust domain.WorkerTrust) error {
+ sql, args, err := psql.Update("workers").
+ SetMap(map[string]any{"trust_level": string(trust), "updated_at": time.Now()}).
+ Where(sq.Eq{"id": id}).
+ ToSql()
+ if err != nil {
+ return err
+ }
+ tag, err := conn(ctx, r.pool).Exec(ctx, sql, args...)
+ if err != nil {
+ return fmt.Errorf("set worker trust: %w", err)
+ }
+ if tag.RowsAffected() == 0 {
+ return domain.ErrWorkerNotFound
+ }
+ return nil
+}
+
func scanWorker(row pgx.Row) (*domain.Worker, error) {
var (
w domain.Worker
diff --git a/coordinator/internal/storage/postgres/workload_settings_repo.go b/coordinator/internal/storage/postgres/workload_settings_repo.go
new file mode 100644
index 0000000..f668f56
--- /dev/null
+++ b/coordinator/internal/storage/postgres/workload_settings_repo.go
@@ -0,0 +1,67 @@
+package postgres
+
+import (
+ "context"
+ "fmt"
+ "time"
+
+ "github.com/jackc/pgx/v5/pgxpool"
+
+ "github.com/emil28092005/SciMesh/coordinator/internal/usecase"
+)
+
+// WorkloadSettingsRepo persists the per-workload enable/disable overrides.
+// Absence of a row means the workload is enabled (the catalog default).
+type WorkloadSettingsRepo struct{ pool *pgxpool.Pool }
+
+func NewWorkloadSettingsRepo(pool *pgxpool.Pool) *WorkloadSettingsRepo {
+ return &WorkloadSettingsRepo{pool: pool}
+}
+
+var _ usecase.WorkloadSettingsRepository = (*WorkloadSettingsRepo)(nil)
+
+func (r *WorkloadSettingsRepo) GetEnabled(ctx context.Context, workload string) (bool, error) {
+ var enabled bool
+ err := conn(ctx, r.pool).QueryRow(ctx,
+ "SELECT enabled FROM workload_settings WHERE workload = $1", workload).Scan(&enabled)
+ if err != nil && err.Error() == "no rows in result set" {
+ return true, nil // no override: catalog default enabled
+ }
+ if err != nil {
+ return false, fmt.Errorf("get workload setting: %w", err)
+ }
+ return enabled, nil
+}
+
+func (r *WorkloadSettingsRepo) List(ctx context.Context) ([]usecase.WorkloadSetting, error) {
+ rows, err := conn(ctx, r.pool).Query(ctx,
+ "SELECT workload, enabled, updated_at FROM workload_settings ORDER BY workload ASC")
+ if err != nil {
+ return nil, fmt.Errorf("list workload settings: %w", err)
+ }
+ defer rows.Close()
+ var out []usecase.WorkloadSetting
+ for rows.Next() {
+ var s usecase.WorkloadSetting
+ if err := rows.Scan(&s.Workload, &s.Enabled, &s.UpdatedAt); err != nil {
+ return nil, err
+ }
+ out = append(out, s)
+ }
+ return out, rows.Err()
+}
+
+func (r *WorkloadSettingsRepo) SetEnabled(ctx context.Context, workload string, enabled bool, now time.Time) error {
+ sql, args, err := psql.Insert("workload_settings").
+ Columns("workload", "enabled", "updated_at").
+ Values(workload, enabled, now).
+ Suffix(`ON CONFLICT (workload) DO UPDATE SET enabled = EXCLUDED.enabled, updated_at = EXCLUDED.updated_at`).
+ ToSql()
+ if err != nil {
+ return err
+ }
+ if _, err := conn(ctx, r.pool).Exec(ctx, sql, args...); err != nil {
+ return fmt.Errorf("set workload setting: %w", err)
+ }
+ return nil
+}
diff --git a/coordinator/internal/storage/sqlite/admin_m2_test.go b/coordinator/internal/storage/sqlite/admin_m2_test.go
new file mode 100644
index 0000000..ef3a060
--- /dev/null
+++ b/coordinator/internal/storage/sqlite/admin_m2_test.go
@@ -0,0 +1,89 @@
+package sqlite
+
+import (
+ "context"
+ "testing"
+ "time"
+
+ "github.com/google/uuid"
+
+ "github.com/emil28092005/SciMesh/coordinator/internal/domain"
+)
+
+func TestWorkloadSettingsRepoRoundTrip(t *testing.T) {
+ db := newTestDB(t)
+ ctx := context.Background()
+ repo := NewWorkloadSettingsRepo(db)
+
+ // No override: enabled by default.
+ enabled, err := repo.GetEnabled(ctx, "similarity-search")
+ if err != nil {
+ t.Fatal(err)
+ }
+ if !enabled {
+ t.Error("workload without an override must be enabled")
+ }
+
+ now := time.Date(2026, 8, 2, 12, 0, 0, 0, time.UTC)
+ if err := repo.SetEnabled(ctx, "similarity-search", false, now); err != nil {
+ t.Fatal(err)
+ }
+ enabled, err = repo.GetEnabled(ctx, "similarity-search")
+ if err != nil {
+ t.Fatal(err)
+ }
+ if enabled {
+ t.Error("workload must be disabled after the override")
+ }
+
+ // Upsert flips it back and updates the timestamp.
+ later := now.Add(time.Hour)
+ if err := repo.SetEnabled(ctx, "similarity-search", true, later); err != nil {
+ t.Fatal(err)
+ }
+ enabled, err = repo.GetEnabled(ctx, "similarity-search")
+ if err != nil {
+ t.Fatal(err)
+ }
+ if !enabled {
+ t.Error("workload must be re-enabled after the upsert")
+ }
+
+ list, err := repo.List(ctx)
+ if err != nil {
+ t.Fatal(err)
+ }
+ if len(list) != 1 || list[0].Workload != "similarity-search" || !list[0].Enabled {
+ t.Errorf("list = %+v, want the single re-enabled override", list)
+ }
+}
+
+func TestWorkerSetTrust(t *testing.T) {
+ db := newTestDB(t)
+ ctx := context.Background()
+ repo := NewWorkerRepo(db)
+
+ worker, err := domain.NewWorker("lab-node", []string{"similarity-search"}, fixedTime())
+ if err != nil {
+ t.Fatal(err)
+ }
+ if err := repo.Insert(ctx, worker); err != nil {
+ t.Fatal(err)
+ }
+ if err := repo.SetTrust(ctx, worker.ID, domain.WorkerUntrusted); err != nil {
+ t.Fatal(err)
+ }
+ got, err := repo.Get(ctx, worker.ID)
+ if err != nil {
+ t.Fatal(err)
+ }
+ if got.TrustLevel != domain.WorkerUntrusted {
+ t.Errorf("trust = %q, want untrusted", got.TrustLevel)
+ }
+ if err := repo.SetTrust(ctx, worker.ID, domain.WorkerTrusted); err != nil {
+ t.Fatal(err)
+ }
+ if err := repo.SetTrust(ctx, uuid.New(), domain.WorkerTrusted); err != domain.ErrWorkerNotFound {
+ t.Errorf("unknown worker trust err = %v, want ErrWorkerNotFound", err)
+ }
+}
diff --git a/coordinator/internal/storage/sqlite/migrations/0002_workload_settings.sql b/coordinator/internal/storage/sqlite/migrations/0002_workload_settings.sql
new file mode 100644
index 0000000..7c3b7a8
--- /dev/null
+++ b/coordinator/internal/storage/sqlite/migrations/0002_workload_settings.sql
@@ -0,0 +1,8 @@
+-- 0002: per-workload enable/disable. Absence of a row means "enabled" (the
+-- catalog default); a row only exists once an admin flipped a workload off or
+-- back on.
+CREATE TABLE IF NOT EXISTS workload_settings (
+ workload TEXT NOT NULL PRIMARY KEY,
+ enabled INTEGER NOT NULL CHECK (enabled IN (0, 1)),
+ updated_at INTEGER NOT NULL
+);
diff --git a/coordinator/internal/storage/sqlite/sqlite_test.go b/coordinator/internal/storage/sqlite/sqlite_test.go
index 78eb679..55792b5 100644
--- a/coordinator/internal/storage/sqlite/sqlite_test.go
+++ b/coordinator/internal/storage/sqlite/sqlite_test.go
@@ -69,8 +69,8 @@ func TestMigrateIsIdempotent(t *testing.T) {
if err := db.QueryRowContext(ctx, "PRAGMA user_version").Scan(&version); err != nil {
t.Fatal(err)
}
- if version != 1 {
- t.Errorf("user_version = %d, want 1", version)
+ if version != 2 {
+ t.Errorf("user_version = %d, want 2", version)
}
}
diff --git a/coordinator/internal/storage/sqlite/worker_repo.go b/coordinator/internal/storage/sqlite/worker_repo.go
index 4d0bf44..3dbfc17 100644
--- a/coordinator/internal/storage/sqlite/worker_repo.go
+++ b/coordinator/internal/storage/sqlite/worker_repo.go
@@ -90,3 +90,20 @@ func (r *WorkerRepo) MarkStaleOffline(ctx context.Context, cutoff time.Time) (in
}
return res.RowsAffected()
}
+
+func (r *WorkerRepo) SetTrust(ctx context.Context, id uuid.UUID, trust domain.WorkerTrust) error {
+ res, err := conn(ctx, r.db).ExecContext(ctx,
+ "UPDATE workers SET trust_level = ?, updated_at = ? WHERE id = ?",
+ string(trust), encodeTime(time.Now()), id.String())
+ if err != nil {
+ return err
+ }
+ affected, err := res.RowsAffected()
+ if err != nil {
+ return err
+ }
+ if affected == 0 {
+ return domain.ErrWorkerNotFound
+ }
+ return nil
+}
diff --git a/coordinator/internal/storage/sqlite/workload_settings_repo.go b/coordinator/internal/storage/sqlite/workload_settings_repo.go
new file mode 100644
index 0000000..c41c53a
--- /dev/null
+++ b/coordinator/internal/storage/sqlite/workload_settings_repo.go
@@ -0,0 +1,75 @@
+package sqlite
+
+import (
+ "context"
+ "database/sql"
+ "fmt"
+ "time"
+
+ "github.com/emil28092005/SciMesh/coordinator/internal/usecase"
+)
+
+// WorkloadSettingsRepo persists the per-workload enable/disable overrides.
+// Absence of a row means the workload is enabled (the catalog default).
+type WorkloadSettingsRepo struct{ db *sql.DB }
+
+func NewWorkloadSettingsRepo(db *sql.DB) *WorkloadSettingsRepo { return &WorkloadSettingsRepo{db: db} }
+
+var _ usecase.WorkloadSettingsRepository = (*WorkloadSettingsRepo)(nil)
+
+func (r *WorkloadSettingsRepo) GetEnabled(ctx context.Context, workload string) (bool, error) {
+ var enabled int
+ err := conn(ctx, r.db).QueryRowContext(ctx,
+ "SELECT enabled FROM workload_settings WHERE workload = ?", workload).Scan(&enabled)
+ if err == sql.ErrNoRows {
+ return true, nil // no override: catalog default enabled
+ }
+ if err != nil {
+ return false, fmt.Errorf("get workload setting: %w", err)
+ }
+ return enabled == 1, nil
+}
+
+func (r *WorkloadSettingsRepo) List(ctx context.Context) ([]usecase.WorkloadSetting, error) {
+ rows, err := conn(ctx, r.db).QueryContext(ctx,
+ "SELECT workload, enabled, updated_at FROM workload_settings ORDER BY workload ASC")
+ if err != nil {
+ return nil, fmt.Errorf("list workload settings: %w", err)
+ }
+ defer func() { _ = rows.Close() }()
+ var out []usecase.WorkloadSetting
+ for rows.Next() {
+ var (
+ name string
+ enabled int
+ updatedAt int64
+ )
+ if err := rows.Scan(&name, &enabled, &updatedAt); err != nil {
+ return nil, err
+ }
+ out = append(out, usecase.WorkloadSetting{
+ Workload: name,
+ Enabled: enabled == 1,
+ UpdatedAt: decodeTime(updatedAt),
+ })
+ }
+ return out, rows.Err()
+}
+
+func (r *WorkloadSettingsRepo) SetEnabled(ctx context.Context, workload string, enabled bool, now time.Time) error {
+ _, err := conn(ctx, r.db).ExecContext(ctx, `
+INSERT INTO workload_settings (workload, enabled, updated_at) VALUES (?, ?, ?)
+ON CONFLICT (workload) DO UPDATE SET enabled = excluded.enabled, updated_at = excluded.updated_at`,
+ workload, boolInt(enabled), now.UnixNano())
+ if err != nil {
+ return fmt.Errorf("set workload setting: %w", err)
+ }
+ return nil
+}
+
+func boolInt(b bool) int {
+ if b {
+ return 1
+ }
+ return 0
+}
diff --git a/coordinator/internal/transport/http/errors.go b/coordinator/internal/transport/http/errors.go
index ae14782..ef8147b 100644
--- a/coordinator/internal/transport/http/errors.go
+++ b/coordinator/internal/transport/http/errors.go
@@ -42,7 +42,7 @@ func (s *Server) writeError(w http.ResponseWriter, r *http.Request, err error) {
status := http.StatusInternalServerError
switch {
- case errors.Is(err, domain.ErrInvalidInput):
+ case errors.Is(err, domain.ErrInvalidInput), errors.Is(err, domain.ErrWorkloadDisabled):
status = http.StatusBadRequest
case errors.Is(err, domain.ErrJobNotFound), errors.Is(err, domain.ErrTaskNotFound),
errors.Is(err, domain.ErrWorkerNotFound), errors.Is(err, domain.ErrArtifactNotFound):
diff --git a/coordinator/internal/transport/http/server.go b/coordinator/internal/transport/http/server.go
index d01edd0..b919846 100644
--- a/coordinator/internal/transport/http/server.go
+++ b/coordinator/internal/transport/http/server.go
@@ -182,6 +182,16 @@ func (s *Server) Handler(token string, uiToken ...string) http.Handler {
ui.Handle("GET /ui/admin/api/system", chain(http.HandlerFunc(s.handleUIAdminSystemJSON), gate, requireAdmin))
ui.Handle("GET /ui/admin/api/jobs", chain(http.HandlerFunc(s.handleUIAdminJobsJSON), gate, requireAdmin))
ui.Handle("GET /ui/admin/api/metrics", chain(http.HandlerFunc(s.handleUIAdminMetricsJSON), gate, requireAdmin))
+ ui.Handle("GET /ui/admin/api/workers", chain(http.HandlerFunc(s.handleUIAdminWorkersJSON), gate, requireAdmin))
+ ui.Handle("POST /ui/admin/api/workers/{id}/trust", chain(http.HandlerFunc(s.handleUIAdminSetTrustJSON), gate, requireAdmin))
+ ui.Handle("GET /ui/admin/api/users", chain(http.HandlerFunc(s.handleUIAdminUsersJSON), gate, requireAdmin))
+ ui.Handle("POST /ui/admin/api/users/{id}/role", chain(http.HandlerFunc(s.handleUIAdminSetUserRoleJSON), gate, requireAdmin))
+ ui.Handle("GET /ui/admin/api/worker-keys", chain(http.HandlerFunc(s.handleUIAdminWorkerKeysJSON), gate, requireAdmin))
+ ui.Handle("POST /ui/admin/api/worker-keys/{id}/revoke", chain(http.HandlerFunc(s.handleUIAdminRevokeKeyJSON), gate, requireAdmin))
+ ui.Handle("GET /ui/admin/api/workloads", chain(http.HandlerFunc(s.handleUIAdminWorkloadsJSON), gate, requireAdmin))
+ ui.Handle("POST /ui/admin/api/workloads/{name}/enabled", chain(http.HandlerFunc(s.handleUIAdminSetWorkloadEnabledJSON), gate, requireAdmin))
+ ui.Handle("GET /ui/admin/api/settings", chain(http.HandlerFunc(s.handleUIAdminSettingsJSON), gate, requireAdmin))
+ ui.Handle("POST /ui/admin/api/token/reveal", chain(http.HandlerFunc(s.handleUIAdminRevealTokenJSON), gate, requireAdmin))
} else {
for _, rt := range app {
ui.HandleFunc(rt.pattern, rt.handler)
diff --git a/coordinator/internal/transport/http/server_test.go b/coordinator/internal/transport/http/server_test.go
index fd6b2f8..3c401b6 100644
--- a/coordinator/internal/transport/http/server_test.go
+++ b/coordinator/internal/transport/http/server_test.go
@@ -39,6 +39,7 @@ func newEnvWithUIToken(t *testing.T, ready func(context.Context) error, configur
work := memstore.NewWorkerRepo()
arts := memstore.NewArtifactRepo()
blobs := memstore.NewBlobStore()
+ settings := memstoreSettings{}
clk := memstore.NewClock(time.Date(2026, 7, 21, 12, 0, 0, 0, time.UTC))
tx := memstore.Tx{}
lease := 2 * time.Minute
@@ -47,7 +48,7 @@ func newEnvWithUIToken(t *testing.T, ready func(context.Context) error, configur
uc := coordhttp.UseCases{
RegisterWorker: usecase.NewRegisterWorker(work, clk),
CreateJob: usecase.NewCreateJob(jobs, tasks, tx, clk),
- SubmitDataset: usecase.NewSubmitDataset(blobs, arts, jobs, tasks, tx, clk, 3, testCatalog()),
+ SubmitDataset: usecase.NewSubmitDataset(blobs, arts, jobs, tasks, tx, clk, 3, testCatalog(), settings),
ClaimTask: usecase.NewClaimTask(tasks, jobs, work, tx, clk, lease, testCatalog()),
RenewLease: usecase.NewRenewLease(tasks, work, tx, clk, lease),
CompleteTask: usecase.NewCompleteTask(tasks, jobs, arts, work, memstore.NewTaskResultRepo(), tx, clk, 2, testCatalog()),
@@ -739,3 +740,22 @@ func (e *env) uploadDataset(t *testing.T, workload string, rows int, tsv string)
}
func itoa(n int) string { return strconv.Itoa(n) }
+
+// memstoreSettings is an in-memory WorkloadSettingsRepository for tests.
+type memstoreSettings struct{ overrides map[string]bool }
+
+func (m memstoreSettings) GetEnabled(ctx context.Context, workload string) (bool, error) {
+ if enabled, ok := m.overrides[workload]; ok {
+ return enabled, nil
+ }
+ return true, nil
+}
+
+func (m memstoreSettings) List(ctx context.Context) ([]usecase.WorkloadSetting, error) {
+ return nil, nil
+}
+
+func (m memstoreSettings) SetEnabled(ctx context.Context, workload string, enabled bool, now time.Time) error {
+ m.overrides[workload] = enabled
+ return nil
+}
diff --git a/coordinator/internal/transport/http/templates/admin.html b/coordinator/internal/transport/http/templates/admin.html
index ba0f603..e920f57 100644
--- a/coordinator/internal/transport/http/templates/admin.html
+++ b/coordinator/internal/transport/http/templates/admin.html
@@ -205,13 +205,35 @@ tbody tr:hover{background:var(--panel-2)}
-
+
-
Workers milestone M2 Trust management lands with milestone M2 Worker list, trust dropdown and heartbeat overview are wired next. The dashboard already shows the live fleet.
+ Workers register themselves. Trust decides whether a machine's results are accepted directly or need quorum.
+
+
+ Worker Status Capabilities Trust Owner Last signal
+
+
+
-
+
+ Users
+
+
+ Email Role Verified Created
+
+
+
+
+ Worker keys
+ Keys let lab machines register as workers under a user account. Served instances can also use the cluster token.
+
+
+ Label Prefix Owner Created Last used
+
+
+
Quick user action
- Accounts and worker keys
- User list and key management land with milestone M2 The userservice gains admin list endpoints; the console gets tables, role selects and key revoke.
-
+
- Workload enable/disable lands with milestone M2 The catalog is already served to the job form; persisted on/off switches arrive with the settings migration.
+ Disabled workloads are rejected at submit time and hidden from the job form. Settings persist in the database.
+
+
+ Workload Reduction Parameters Dataset upload Enabled
+
+
+
@@ -255,9 +281,20 @@ tbody tr:hover{background:var(--panel-2)}
-
+
- Cluster settings land with milestone M2 Worker token reveal, public URL and storage settings follow in the next milestone.
+ The cluster token below authenticates any worker. Reveal it only on a trusted machine.
+ Cluster
+
+
+ Worker token ••••••••••••••••••••••••Reveal
+ Public URL —
+ Listen address —
+ Data directory —
+ Database engine —
+ Binary —
+
+
@@ -272,10 +309,10 @@ const fmtBytes=b=>{if(b==null||b<0)return '—';if(b<1024)return b+' B';if(b<104
const fmtTime=t=>t?new Date(t).toLocaleString():'—';
const esc=s=>String(s).replace(/[&<>"]/g,c=>({'&':'&','<':'<','>':'>','"':'"'}[c]));
let timer=null,current={page:'system',jobsStatus:'',jobsPage:1};
-const setPage=page=>{current.page=page;document.querySelectorAll('.nav-item').forEach(b=>b.classList.toggle('active',b.dataset.page===page));document.querySelectorAll('.page').forEach(p=>p.classList.toggle('active',p.id==='page-'+page));const [t,s]=titles[page];document.getElementById('page-title').textContent=t;document.getElementById('page-sub').textContent=s;if(timer){clearInterval(timer);timer=null}if(page==='system'||page==='jobs'||page==='metrics'){refresh();timer=setInterval(refresh,5000)}};
+const setPage=page=>{current.page=page;document.querySelectorAll('.nav-item').forEach(b=>b.classList.toggle('active',b.dataset.page===page));document.querySelectorAll('.page').forEach(p=>p.classList.toggle('active',p.id==='page-'+page));const [t,s]=titles[page];document.getElementById('page-title').textContent=t;document.getElementById('page-sub').textContent=s;if(timer){clearInterval(timer);timer=null}refresh();timer=setInterval(refresh,5000)};
document.querySelectorAll('.nav-item').forEach(b=>b.addEventListener('click',()=>setPage(b.dataset.page)));
-const refresh=()=>{if(document.hidden)return;const p=current.page;if(p==='system')loadSystem();else if(p==='jobs')loadJobs();else if(p==='metrics')loadMetrics()};
+const refresh=()=>{if(document.hidden)return;const p=current.page;if(p==='system')loadSystem();else if(p==='jobs')loadJobs();else if(p==='metrics')loadMetrics();else if(p==='workers')loadWorkers();else if(p==='users')loadUsers();else if(p==='workloads')loadWorkloads();else if(p==='settings')loadSettings()};
document.addEventListener('visibilitychange',()=>{if(!document.hidden)refresh()});
async function loadSystem(){
@@ -384,6 +421,114 @@ async function loadMetrics(){
}
}
setPage('system');
+
+const workerStatusPill=s=>({online:['Online','pill-success'],busy:['Busy','pill-active'],offline:['Offline','pill-waiting']}[s]||[s,'pill-waiting']);
+async function loadWorkers(){
+ const r=await fetch('/ui/admin/api/workers',{headers:{Accept:'application/json'}});
+ if(!r.ok)return;
+ const v=await r.json();
+ const rows=document.getElementById('worker-rows');
+ rows.replaceChildren();
+ if(!v.workers.length){const tr=document.createElement('tr');tr.innerHTML='
No worker is registered yet.
';rows.append(tr);return}
+ for(const w of v.workers){
+ const tr=document.createElement('tr');
+ const [label,cls]=workerStatusPill(w.status);
+ const trustSel='Trusted Untrusted ';
+ tr.innerHTML=''+esc(w.name)+'
'+esc(w.id.slice(0,8))+'…
'+
+ ''+pill(label,cls,null)+' '+
+ ''+(w.capabilities||[]).map(c=>''+esc(c)+' ').join('')+' '+
+ ''+trustSel+' '+
+ ''+esc(w.owner)+' '+
+ ''+fmtTime(w.last_heartbeat_at)+' ';
+ rows.append(tr);
+ }
+ document.querySelectorAll('.trust-sel').forEach(sel=>sel.addEventListener('change',async()=>{
+ await fetch('/ui/admin/api/workers/'+sel.dataset.id+'/trust',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({trusted:sel.value==='trusted'})});
+ loadWorkers();
+ }));
+}
+async function loadUsers(){
+ const r=await fetch('/ui/admin/api/users',{headers:{Accept:'application/json'}});
+ if(!r.ok)return;
+ const v=await r.json();
+ const rows=document.getElementById('user-rows');
+ rows.replaceChildren();
+ for(const u of v.users||[]){
+ const tr=document.createElement('tr');
+ const roleSel='user admin ';
+ tr.innerHTML=''+esc(u.email)+'
'+
+ ''+roleSel+' '+
+ ''+pill(u.verified?'Verified':'—',u.verified?'pill-success':'pill-waiting',null)+' '+
+ ''+fmtTime(u.created_at)+' ';
+ rows.append(tr);
+ }
+ document.getElementById('user-count').textContent=(v.users||[]).length+' users';
+ document.querySelectorAll('.role-sel').forEach(sel=>sel.addEventListener('change',async()=>{
+ await fetch('/ui/admin/api/users/'+sel.dataset.id+'/role',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({role:sel.value})});
+ loadUsers();
+ }));
+ const keys=await (await fetch('/ui/admin/api/worker-keys',{headers:{Accept:'application/json'}})).json();
+ const keyRows=document.getElementById('key-rows');
+ keyRows.replaceChildren();
+ const emailOf={};for(const u of v.users||[])emailOf[u.id]=u.email;
+ for(const k of keys.worker_keys||[]){
+ const tr=document.createElement('tr');
+ tr.innerHTML=''+esc(k.name)+' '+
+ ''+esc(k.prefix)+'… '+
+ ''+esc(emailOf[k.user_id]||k.user_id.slice(0,8)+'…')+' '+
+ ''+fmtTime(k.created_at)+' '+
+ ''+(k.last_used_at?fmtTime(k.last_used_at):'never')+' '+
+ ''+((k.revoked_at)?' Revoked ':'Revoke ')+' ';
+ keyRows.append(tr);
+ }
+ document.querySelectorAll('.key-revoke').forEach(btn=>btn.addEventListener('click',async()=>{
+ if(!confirm('Revoke this worker key? The machine will be cut off on its next refresh.'))return;
+ await fetch('/ui/admin/api/worker-keys/'+btn.dataset.id+'/revoke',{method:'POST'});
+ loadUsers();
+ }));
+}
+async function loadWorkloads(){
+ const r=await fetch('/ui/admin/api/workloads',{headers:{Accept:'application/json'}});
+ if(!r.ok)return;
+ const v=await r.json();
+ const rows=document.getElementById('workload-rows');
+ rows.replaceChildren();
+ for(const w of v.workloads||[]){
+ const tr=document.createElement('tr');
+ tr.innerHTML=''+esc(w.name)+'
'+esc(w.description||'')+'
'+
+ ' '+esc(w.reduction)+' '+
+ ''+w.parameters+' declared '+
+ ''+pill(w.upload_ready?'ready':'—',w.upload_ready?'pill-success':'pill-waiting',null)+' '+
+ ' ';
+ rows.append(tr);
+ }
+ document.querySelectorAll('.wl-toggle').forEach(t=>t.addEventListener('click',async()=>{
+ const enabled=!t.classList.contains('on');
+ await fetch('/ui/admin/api/workloads/'+encodeURIComponent(t.dataset.name)+'/enabled',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({enabled})});
+ t.classList.toggle('on',enabled);
+ }));
+}
+async function loadSettings(){
+ const r=await fetch('/ui/admin/api/settings',{headers:{Accept:'application/json'}});
+ if(!r.ok)return;
+ const v=await r.json();
+ document.getElementById('s-public').textContent=v.public_url||'—';
+ document.getElementById('s-addr').textContent=v.addr;
+ document.getElementById('s-datadir').textContent=v.data_dir||'—';
+ document.getElementById('s-engine').textContent=v.db_engine;
+ document.getElementById('s-binary').textContent=v.binary||'—';
+}
+document.getElementById('reveal').addEventListener('click',async e=>{
+ const tok=document.getElementById('tok');
+ if(tok.textContent.startsWith('•')){
+ const r=await fetch('/ui/admin/api/token/reveal',{method:'POST',headers:{Accept:'application/json'}});
+ if(!r.ok)return;
+ const v=await r.json();
+ tok.textContent=v.token||'(none)';
+ e.target.textContent='Hide';
+ }else{tok.textContent='••••••••••••••••••••••••';e.target.textContent='Reveal'}
+});
+document.querySelectorAll('.toggle').forEach(t=>t.addEventListener('click',()=>t.classList.toggle('on')));