feat(service): add TaskService, FileStagingService, and refactor ApplicationService for task submission

This commit is contained in:
dailz
2026-04-15 21:31:02 +08:00
parent acf8c1d62b
commit ec64300ff2
9 changed files with 2394 additions and 136 deletions

View File

@@ -5,6 +5,8 @@ import (
"encoding/json"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strings"
"testing"
@@ -42,24 +44,22 @@ func setupApplicationService(t *testing.T, slurmHandler http.HandlerFunc) (*Appl
}
func TestValidateParams_AllRequired(t *testing.T) {
svc := NewApplicationService(nil, nil, "", zap.NewNop())
params := []model.ParameterSchema{
{Name: "NAME", Type: model.ParamTypeString, Required: true},
{Name: "COUNT", Type: model.ParamTypeInteger, Required: true},
}
values := map[string]string{"NAME": "hello", "COUNT": "5"}
if err := svc.ValidateParams(params, values); err != nil {
if err := ValidateParams(params, values); err != nil {
t.Errorf("expected no error, got %v", err)
}
}
func TestValidateParams_MissingRequired(t *testing.T) {
svc := NewApplicationService(nil, nil, "", zap.NewNop())
params := []model.ParameterSchema{
{Name: "NAME", Type: model.ParamTypeString, Required: true},
}
values := map[string]string{}
err := svc.ValidateParams(params, values)
err := ValidateParams(params, values)
if err == nil {
t.Fatal("expected error for missing required param")
}
@@ -69,12 +69,11 @@ func TestValidateParams_MissingRequired(t *testing.T) {
}
func TestValidateParams_InvalidInteger(t *testing.T) {
svc := NewApplicationService(nil, nil, "", zap.NewNop())
params := []model.ParameterSchema{
{Name: "COUNT", Type: model.ParamTypeInteger, Required: true},
}
values := map[string]string{"COUNT": "abc"}
err := svc.ValidateParams(params, values)
err := ValidateParams(params, values)
if err == nil {
t.Fatal("expected error for invalid integer")
}
@@ -84,12 +83,11 @@ func TestValidateParams_InvalidInteger(t *testing.T) {
}
func TestValidateParams_InvalidEnum(t *testing.T) {
svc := NewApplicationService(nil, nil, "", zap.NewNop())
params := []model.ParameterSchema{
{Name: "MODE", Type: model.ParamTypeEnum, Required: true, Options: []string{"fast", "slow"}},
}
values := map[string]string{"MODE": "medium"}
err := svc.ValidateParams(params, values)
err := ValidateParams(params, values)
if err == nil {
t.Fatal("expected error for invalid enum value")
}
@@ -99,12 +97,11 @@ func TestValidateParams_InvalidEnum(t *testing.T) {
}
func TestValidateParams_BooleanValues(t *testing.T) {
svc := NewApplicationService(nil, nil, "", zap.NewNop())
params := []model.ParameterSchema{
{Name: "FLAG", Type: model.ParamTypeBoolean, Required: true},
}
for _, val := range []string{"true", "false", "1", "0"} {
err := svc.ValidateParams(params, map[string]string{"FLAG": val})
err := ValidateParams(params, map[string]string{"FLAG": val})
if err != nil {
t.Errorf("boolean value %q should be valid, got error: %v", val, err)
}
@@ -112,10 +109,9 @@ func TestValidateParams_BooleanValues(t *testing.T) {
}
func TestRenderScript_SimpleReplacement(t *testing.T) {
svc := NewApplicationService(nil, nil, "", zap.NewNop())
params := []model.ParameterSchema{{Name: "INPUT", Type: model.ParamTypeString}}
values := map[string]string{"INPUT": "data.txt"}
result := svc.RenderScript("echo $INPUT", params, values)
result := RenderScript("echo $INPUT", params, values)
expected := "echo 'data.txt'"
if result != expected {
t.Errorf("got %q, want %q", result, expected)
@@ -123,10 +119,9 @@ func TestRenderScript_SimpleReplacement(t *testing.T) {
}
func TestRenderScript_DefaultValues(t *testing.T) {
svc := NewApplicationService(nil, nil, "", zap.NewNop())
params := []model.ParameterSchema{{Name: "OUTPUT", Type: model.ParamTypeString, Default: "out.log"}}
values := map[string]string{}
result := svc.RenderScript("cat $OUTPUT", params, values)
result := RenderScript("cat $OUTPUT", params, values)
expected := "cat 'out.log'"
if result != expected {
t.Errorf("got %q, want %q", result, expected)
@@ -134,10 +129,9 @@ func TestRenderScript_DefaultValues(t *testing.T) {
}
func TestRenderScript_PreservesUnknownVars(t *testing.T) {
svc := NewApplicationService(nil, nil, "", zap.NewNop())
params := []model.ParameterSchema{{Name: "INPUT", Type: model.ParamTypeString}}
values := map[string]string{"INPUT": "data.txt"}
result := svc.RenderScript("export HOME=$HOME\necho $INPUT\necho $PATH", params, values)
result := RenderScript("export HOME=$HOME\necho $INPUT\necho $PATH", params, values)
if !strings.Contains(result, "$HOME") {
t.Error("$HOME should be preserved")
}
@@ -150,7 +144,6 @@ func TestRenderScript_PreservesUnknownVars(t *testing.T) {
}
func TestRenderScript_ShellEscaping(t *testing.T) {
svc := NewApplicationService(nil, nil, "", zap.NewNop())
params := []model.ParameterSchema{{Name: "INPUT", Type: model.ParamTypeString}}
tests := []struct {
@@ -165,7 +158,7 @@ func TestRenderScript_ShellEscaping(t *testing.T) {
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
result := svc.RenderScript("$INPUT", params, map[string]string{"INPUT": tt.value})
result := RenderScript("$INPUT", params, map[string]string{"INPUT": tt.value})
if result != tt.expected {
t.Errorf("got %q, want %q", result, tt.expected)
}
@@ -174,14 +167,13 @@ func TestRenderScript_ShellEscaping(t *testing.T) {
}
func TestRenderScript_OverlappingParams(t *testing.T) {
svc := NewApplicationService(nil, nil, "", zap.NewNop())
template := "$JOB_NAME and $JOB"
params := []model.ParameterSchema{
{Name: "JOB", Type: model.ParamTypeString},
{Name: "JOB_NAME", Type: model.ParamTypeString},
}
values := map[string]string{"JOB": "myjob", "JOB_NAME": "my-test-job"}
result := svc.RenderScript(template, params, values)
result := RenderScript(template, params, values)
if strings.Contains(result, "$JOB_NAME") {
t.Error("$JOB_NAME was not replaced")
}
@@ -197,10 +189,9 @@ func TestRenderScript_OverlappingParams(t *testing.T) {
}
func TestRenderScript_NewlineInValue(t *testing.T) {
svc := NewApplicationService(nil, nil, "", zap.NewNop())
params := []model.ParameterSchema{{Name: "CMD", Type: model.ParamTypeString}}
values := map[string]string{"CMD": "line1\nline2"}
result := svc.RenderScript("echo $CMD", params, values)
result := RenderScript("echo $CMD", params, values)
expected := "echo 'line1\nline2'"
if result != expected {
t.Errorf("got %q, want %q", result, expected)
@@ -298,3 +289,68 @@ func TestSubmitFromApplication_NoParameters(t *testing.T) {
t.Errorf("JobID = %d, want 99", resp.JobID)
}
}
func TestSubmitFromApplication_DelegatesToTaskService(t *testing.T) {
jobID := int32(77)
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
json.NewEncoder(w).Encode(slurm.OpenapiJobSubmitResponse{
Result: &slurm.JobSubmitResponseMsg{JobID: &jobID},
})
}))
defer srv.Close()
client, _ := slurm.NewClient(srv.URL, srv.Client())
jobSvc := NewJobService(client, zap.NewNop())
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{
Logger: gormlogger.Default.LogMode(gormlogger.Silent),
})
if err != nil {
t.Fatalf("open sqlite: %v", err)
}
if err := db.AutoMigrate(&model.Application{}, &model.Task{}, &model.File{}, &model.FileBlob{}); err != nil {
t.Fatalf("auto migrate: %v", err)
}
appStore := store.NewApplicationStore(db)
taskStore := store.NewTaskStore(db)
fileStore := store.NewFileStore(db)
blobStore := store.NewBlobStore(db)
workDirBase := filepath.Join(t.TempDir(), "workdir")
os.MkdirAll(workDirBase, 0777)
taskSvc := NewTaskService(taskStore, appStore, fileStore, blobStore, nil, jobSvc, workDirBase, zap.NewNop())
appSvc := NewApplicationService(appStore, jobSvc, workDirBase, zap.NewNop(), taskSvc)
id, err := appStore.Create(context.Background(), &model.CreateApplicationRequest{
Name: "delegated-app",
ScriptTemplate: "#!/bin/bash\n#SBATCH --job-name=$JOB_NAME\necho $INPUT",
Parameters: json.RawMessage(`[{"name":"JOB_NAME","type":"string","required":true},{"name":"INPUT","type":"string","required":true}]`),
})
if err != nil {
t.Fatalf("create app: %v", err)
}
resp, err := appSvc.SubmitFromApplication(context.Background(), id, map[string]string{
"JOB_NAME": "delegated-job",
"INPUT": "test-data",
})
if err != nil {
t.Fatalf("SubmitFromApplication() error = %v", err)
}
if resp.JobID != 77 {
t.Errorf("JobID = %d, want 77", resp.JobID)
}
var task model.Task
if err := db.Where("app_id = ?", id).First(&task).Error; err != nil {
t.Fatalf("no hpc_tasks record found for app_id %d: %v", id, err)
}
if task.SlurmJobID == nil || *task.SlurmJobID != 77 {
t.Errorf("task SlurmJobID = %v, want 77", task.SlurmJobID)
}
if task.Status != model.TaskStatusQueued {
t.Errorf("task Status = %q, want %q", task.Status, model.TaskStatusQueued)
}
}