414 lines
13 KiB
Go
414 lines
13 KiB
Go
package webui
|
||
|
||
import (
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"io"
|
||
"net/http"
|
||
"os"
|
||
"strings"
|
||
"time"
|
||
|
||
"bee/audit/internal/platform"
|
||
)
|
||
|
||
func (h *handler) handleAPIAuditRun(w http.ResponseWriter, _ *http.Request) {
|
||
if h.opts.App == nil {
|
||
writeError(w, http.StatusServiceUnavailable, "app not configured")
|
||
return
|
||
}
|
||
t := &Task{
|
||
ID: newJobID("audit"),
|
||
Name: "Audit",
|
||
Target: "audit",
|
||
Priority: defaultTaskPriority("audit", taskParams{}),
|
||
Status: TaskPending,
|
||
CreatedAt: time.Now(),
|
||
}
|
||
globalQueue.enqueue(t)
|
||
writeJSON(w, map[string]string{"task_id": t.ID, "job_id": t.ID})
|
||
}
|
||
|
||
func (h *handler) handleAPIAuditStream(w http.ResponseWriter, r *http.Request) {
|
||
id := r.URL.Query().Get("job_id")
|
||
if id == "" {
|
||
id = r.URL.Query().Get("task_id")
|
||
}
|
||
if j, ok := globalQueue.findJob(id); ok {
|
||
streamJob(w, r, j)
|
||
return
|
||
}
|
||
http.Error(w, "job not found", http.StatusNotFound)
|
||
}
|
||
|
||
// ── SAT ───────────────────────────────────────────────────────────────────────
|
||
|
||
func (h *handler) handleAPISATRun(target string) http.HandlerFunc {
|
||
return func(w http.ResponseWriter, r *http.Request) {
|
||
if h.opts.App == nil {
|
||
writeError(w, http.StatusServiceUnavailable, "app not configured")
|
||
return
|
||
}
|
||
|
||
var body struct {
|
||
Duration int `json:"duration"`
|
||
StressMode bool `json:"stress_mode"`
|
||
GPUIndices []int `json:"gpu_indices"`
|
||
ExcludeGPUIndices []int `json:"exclude_gpu_indices"`
|
||
StaggerGPUStart bool `json:"stagger_gpu_start"`
|
||
ParallelGPUs bool `json:"parallel_gpus"`
|
||
Loader string `json:"loader"`
|
||
Profile string `json:"profile"`
|
||
DisplayName string `json:"display_name"`
|
||
PlatformComponents []string `json:"platform_components"`
|
||
}
|
||
if r.Body != nil {
|
||
if err := json.NewDecoder(r.Body).Decode(&body); err != nil && !errors.Is(err, io.EOF) {
|
||
writeError(w, http.StatusBadRequest, "invalid request body")
|
||
return
|
||
}
|
||
}
|
||
|
||
params := taskParams{
|
||
Duration: body.Duration,
|
||
StressMode: body.StressMode,
|
||
GPUIndices: body.GPUIndices,
|
||
ExcludeGPUIndices: body.ExcludeGPUIndices,
|
||
StaggerGPUStart: body.StaggerGPUStart,
|
||
ParallelGPUs: body.ParallelGPUs,
|
||
Loader: body.Loader,
|
||
BurnProfile: body.Profile,
|
||
DisplayName: body.DisplayName,
|
||
PlatformComponents: body.PlatformComponents,
|
||
}
|
||
tasks, err := h.enqueueSATTarget(target, params)
|
||
if err != nil {
|
||
writeError(w, http.StatusBadRequest, err.Error())
|
||
return
|
||
}
|
||
writeTaskRunResponse(w, tasks)
|
||
}
|
||
}
|
||
|
||
// enqueueSATTarget builds the task set for one SAT target (splitting
|
||
// homogeneous multi-GPU NVIDIA targets as needed) and enqueues it. Shared by
|
||
// the single-target /api/sat/<target>/run endpoints and /api/sat/run-all.
|
||
func (h *handler) enqueueSATTarget(target string, params taskParams) ([]*Task, error) {
|
||
name := taskDisplayName(target, params.BurnProfile, params.Loader)
|
||
if strings.TrimSpace(params.DisplayName) != "" {
|
||
name = params.DisplayName
|
||
}
|
||
tasks, err := buildNvidiaTaskSet(target, defaultTaskPriority(target, params), time.Now(), params, name, h.opts.App, "sat-"+target)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
for _, t := range tasks {
|
||
globalQueue.enqueue(t)
|
||
}
|
||
return tasks, nil
|
||
}
|
||
|
||
// ── Scenario ─────────────────────────────────────────────────────────────────
|
||
|
||
func (h *handler) handleAPIScenarioList(w http.ResponseWriter, _ *http.Request) {
|
||
if h.opts.App == nil {
|
||
writeError(w, http.StatusServiceUnavailable, "app not configured")
|
||
return
|
||
}
|
||
files, err := h.opts.App.ListAvailableScenarios()
|
||
if err != nil {
|
||
writeError(w, http.StatusInternalServerError, err.Error())
|
||
return
|
||
}
|
||
type scenarioFile struct {
|
||
Name string `json:"name"`
|
||
Description string `json:"description"`
|
||
Device string `json:"device"`
|
||
}
|
||
out := make([]scenarioFile, 0, len(files))
|
||
for _, f := range files {
|
||
out = append(out, scenarioFile{Name: f.Name, Description: f.Description, Device: f.Device})
|
||
}
|
||
writeJSON(w, out)
|
||
}
|
||
|
||
func (h *handler) handleAPIScenarioRun(w http.ResponseWriter, r *http.Request) {
|
||
if h.opts.App == nil {
|
||
writeError(w, http.StatusServiceUnavailable, "app not configured")
|
||
return
|
||
}
|
||
var body struct {
|
||
Name string `json:"name"`
|
||
}
|
||
if r.Body != nil {
|
||
if err := json.NewDecoder(r.Body).Decode(&body); err != nil && !errors.Is(err, io.EOF) {
|
||
writeError(w, http.StatusBadRequest, "invalid request body")
|
||
return
|
||
}
|
||
}
|
||
name := strings.TrimSpace(body.Name)
|
||
if name == "" {
|
||
writeError(w, http.StatusBadRequest, "scenario name is required")
|
||
return
|
||
}
|
||
t := &Task{
|
||
ID: newJobID("scenario"),
|
||
Name: "Scenario: " + name,
|
||
Target: "scenario",
|
||
Priority: defaultTaskPriority("scenario", taskParams{}),
|
||
Status: TaskPending,
|
||
CreatedAt: time.Now(),
|
||
}
|
||
t.params.ScenarioName = name
|
||
globalQueue.enqueue(t)
|
||
writeJSON(w, map[string]string{"task_id": t.ID, "job_id": t.ID})
|
||
}
|
||
|
||
func (h *handler) handleAPIBenchmarkNvidiaRunKind(target string) http.HandlerFunc {
|
||
return func(w http.ResponseWriter, r *http.Request) {
|
||
if h.opts.App == nil {
|
||
writeError(w, http.StatusServiceUnavailable, "app not configured")
|
||
return
|
||
}
|
||
|
||
var body struct {
|
||
Profile string `json:"profile"`
|
||
SizeMB int `json:"size_mb"`
|
||
GPUIndices []int `json:"gpu_indices"`
|
||
ExcludeGPUIndices []int `json:"exclude_gpu_indices"`
|
||
RunNCCL *bool `json:"run_nccl"`
|
||
ParallelGPUs *bool `json:"parallel_gpus"`
|
||
RampUp *bool `json:"ramp_up"`
|
||
DisplayName string `json:"display_name"`
|
||
}
|
||
if r.Body != nil {
|
||
if err := json.NewDecoder(r.Body).Decode(&body); err != nil && !errors.Is(err, io.EOF) {
|
||
writeError(w, http.StatusBadRequest, "invalid request body")
|
||
return
|
||
}
|
||
}
|
||
|
||
runNCCL := true
|
||
if body.RunNCCL != nil {
|
||
runNCCL = *body.RunNCCL
|
||
}
|
||
parallelGPUs := false
|
||
if body.ParallelGPUs != nil {
|
||
parallelGPUs = *body.ParallelGPUs
|
||
}
|
||
rampUp := false
|
||
if body.RampUp != nil {
|
||
rampUp = *body.RampUp
|
||
}
|
||
// Build a descriptive base name that includes profile and mode so the task
|
||
// list is self-explanatory without opening individual task detail pages.
|
||
profile := strings.TrimSpace(body.Profile)
|
||
if profile == "" {
|
||
profile = "standard"
|
||
}
|
||
name := taskDisplayName(target, "", "")
|
||
if strings.TrimSpace(body.DisplayName) != "" {
|
||
name = body.DisplayName
|
||
}
|
||
// Append profile tag.
|
||
name = fmt.Sprintf("%s · %s", name, profile)
|
||
|
||
if target == "nvidia-bench-power" && parallelGPUs {
|
||
writeError(w, http.StatusBadRequest, "power / thermal fit benchmark uses sequential or ramp-up modes only")
|
||
return
|
||
}
|
||
|
||
if rampUp && len(body.GPUIndices) > 1 {
|
||
// Ramp-up mode: RunNvidiaPowerBench internally ramps from 1 to N GPUs
|
||
// in Phase 2 (one additional GPU per step). A single task with all
|
||
// selected GPUs is sufficient — spawning N tasks with growing subsets
|
||
// would repeat all earlier steps redundantly.
|
||
gpus, err := apiListNvidiaGPUs(h.opts.App)
|
||
if err != nil {
|
||
writeError(w, http.StatusBadRequest, err.Error())
|
||
return
|
||
}
|
||
resolved, err := expandSelectedGPUIndices(gpus, body.GPUIndices, body.ExcludeGPUIndices)
|
||
if err != nil {
|
||
writeError(w, http.StatusBadRequest, err.Error())
|
||
return
|
||
}
|
||
if len(resolved) < 2 {
|
||
// Fall through to normal single-task path.
|
||
rampUp = false
|
||
} else {
|
||
now := time.Now()
|
||
rampRunID := fmt.Sprintf("ramp-%s", now.UTC().Format("20060102-150405"))
|
||
taskName := fmt.Sprintf("%s · ramp 1–%d · GPU %s", name, len(resolved), formatGPUIndexList(resolved))
|
||
t := &Task{
|
||
ID: newJobID("bee-bench-nvidia"),
|
||
Name: taskName,
|
||
Target: target,
|
||
Priority: defaultTaskPriority(target, taskParams{}),
|
||
Status: TaskPending,
|
||
CreatedAt: now,
|
||
params: taskParams{
|
||
GPUIndices: append([]int(nil), resolved...),
|
||
SizeMB: body.SizeMB,
|
||
BenchmarkProfile: body.Profile,
|
||
RunNCCL: runNCCL,
|
||
ParallelGPUs: true,
|
||
RampTotal: len(resolved),
|
||
RampRunID: rampRunID,
|
||
DisplayName: taskName,
|
||
},
|
||
}
|
||
globalQueue.enqueue(t)
|
||
writeTaskRunResponse(w, []*Task{t})
|
||
return
|
||
}
|
||
}
|
||
|
||
// For non-ramp tasks append mode tag.
|
||
if parallelGPUs {
|
||
name = fmt.Sprintf("%s · parallel", name)
|
||
} else {
|
||
name = fmt.Sprintf("%s · sequential", name)
|
||
}
|
||
|
||
params := taskParams{
|
||
GPUIndices: body.GPUIndices,
|
||
ExcludeGPUIndices: body.ExcludeGPUIndices,
|
||
SizeMB: body.SizeMB,
|
||
BenchmarkProfile: body.Profile,
|
||
RunNCCL: runNCCL,
|
||
ParallelGPUs: parallelGPUs,
|
||
DisplayName: body.DisplayName,
|
||
}
|
||
tasks, err := buildNvidiaTaskSet(target, defaultTaskPriority(target, params), time.Now(), params, name, h.opts.App, "bee-bench-nvidia")
|
||
if err != nil {
|
||
writeError(w, http.StatusBadRequest, err.Error())
|
||
return
|
||
}
|
||
for _, t := range tasks {
|
||
globalQueue.enqueue(t)
|
||
}
|
||
writeTaskRunResponse(w, tasks)
|
||
}
|
||
}
|
||
|
||
func (h *handler) handleAPIBenchmarkAutotuneRun() http.HandlerFunc {
|
||
return func(w http.ResponseWriter, r *http.Request) {
|
||
if h.opts.App == nil {
|
||
writeError(w, http.StatusServiceUnavailable, "app not configured")
|
||
return
|
||
}
|
||
var body struct {
|
||
Profile string `json:"profile"`
|
||
BenchmarkKind string `json:"benchmark_kind"`
|
||
SizeMB int `json:"size_mb"`
|
||
}
|
||
if r.Body != nil {
|
||
if err := json.NewDecoder(r.Body).Decode(&body); err != nil && !errors.Is(err, io.EOF) {
|
||
writeError(w, http.StatusBadRequest, "invalid request body")
|
||
return
|
||
}
|
||
}
|
||
profile := strings.TrimSpace(body.Profile)
|
||
if profile == "" {
|
||
profile = "standard"
|
||
}
|
||
benchmarkKind := strings.TrimSpace(body.BenchmarkKind)
|
||
if benchmarkKind == "" {
|
||
benchmarkKind = "power-fit"
|
||
}
|
||
now := time.Now()
|
||
taskName := fmt.Sprintf("NVIDIA Benchmark Autotune · %s · %s", profile, benchmarkKind)
|
||
t := &Task{
|
||
ID: newJobID("bee-bench-autotune"),
|
||
Name: taskName,
|
||
Target: "nvidia-bench-autotune",
|
||
Priority: defaultTaskPriority("nvidia-bench-autotune", taskParams{}),
|
||
Status: TaskPending,
|
||
CreatedAt: now,
|
||
params: taskParams{
|
||
BenchmarkProfile: profile,
|
||
BenchmarkKind: benchmarkKind,
|
||
SizeMB: body.SizeMB,
|
||
DisplayName: taskName,
|
||
},
|
||
}
|
||
globalQueue.enqueue(t)
|
||
writeTaskRunResponse(w, []*Task{t})
|
||
}
|
||
}
|
||
|
||
func (h *handler) handleAPIBenchmarkAutotuneStatus(w http.ResponseWriter, r *http.Request) {
|
||
if h.opts.App == nil {
|
||
writeError(w, http.StatusServiceUnavailable, "app not configured")
|
||
return
|
||
}
|
||
cfg, err := h.opts.App.LoadBenchmarkPowerAutotune()
|
||
if err != nil {
|
||
if os.IsNotExist(err) {
|
||
w.WriteHeader(http.StatusOK)
|
||
writeJSON(w, map[string]any{
|
||
"configured": false,
|
||
"decision": platform.ResolveSystemPowerDecision(h.opts.ExportDir),
|
||
})
|
||
return
|
||
}
|
||
writeError(w, http.StatusInternalServerError, err.Error())
|
||
return
|
||
}
|
||
w.WriteHeader(http.StatusOK)
|
||
writeJSON(w, map[string]any{
|
||
"configured": true,
|
||
"config": cfg,
|
||
"decision": platform.ResolveSystemPowerDecision(h.opts.ExportDir),
|
||
})
|
||
}
|
||
|
||
func (h *handler) handleAPIBenchmarkNvidiaRun(w http.ResponseWriter, r *http.Request) {
|
||
h.handleAPIBenchmarkNvidiaRunKind("nvidia-bench-perf").ServeHTTP(w, r)
|
||
}
|
||
|
||
func (h *handler) handleAPISATStream(w http.ResponseWriter, r *http.Request) {
|
||
id := r.URL.Query().Get("job_id")
|
||
if id == "" {
|
||
id = r.URL.Query().Get("task_id")
|
||
}
|
||
if j, ok := globalQueue.findJob(id); ok {
|
||
streamJob(w, r, j)
|
||
return
|
||
}
|
||
http.Error(w, "job not found", http.StatusNotFound)
|
||
}
|
||
|
||
func (h *handler) handleAPISATAbort(w http.ResponseWriter, r *http.Request) {
|
||
id := r.URL.Query().Get("job_id")
|
||
if id == "" {
|
||
id = r.URL.Query().Get("task_id")
|
||
}
|
||
if t, ok := globalQueue.findByID(id); ok {
|
||
globalQueue.mu.Lock()
|
||
switch t.Status {
|
||
case TaskPending:
|
||
t.Status = TaskCancelled
|
||
now := time.Now()
|
||
t.DoneAt = &now
|
||
case TaskRunning:
|
||
if t.job == nil || !t.job.abort() {
|
||
globalQueue.mu.Unlock()
|
||
writeJSON(w, map[string]string{"status": "not_running"})
|
||
return
|
||
}
|
||
globalQueue.mu.Unlock()
|
||
writeJSON(w, map[string]string{"status": "aborting"})
|
||
return
|
||
}
|
||
globalQueue.mu.Unlock()
|
||
writeJSON(w, map[string]string{"status": "aborted"})
|
||
return
|
||
}
|
||
http.Error(w, "job not found", http.StatusNotFound)
|
||
}
|
||
|
||
// ── Services ──────────────────────────────────────────────────────────────────
|