package spindle import ( "context" "io" "log/slog" "net/http" "net/http/httptest" "path/filepath" "testing" "time" kgit "tangled.org/core/knotserver/git" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/mill/executor" "tangled.org/core/spindle/models" "tangled.org/core/spindle/observability" ) func TestHasSkipCIPushOption(t *testing.T) { tests := []struct { name string pushOptions []string want bool }{ { name: "skip-ci requests skip", pushOptions: []string{"skip-ci"}, want: true, }, { name: "ci-skip requests skip", pushOptions: []string{"ci-skip"}, want: true, }, { name: "unrelated ci options do not skip", pushOptions: []string{"verbose-ci", "ci-verbose"}, want: false, }, { name: "empty options do not skip", pushOptions: []string{}, want: false, }, { name: "nil options do not skip", pushOptions: nil, want: false, }, { name: "mixed options skip when any skip option appears", pushOptions: []string{"verbose-ci", "skip-ci", "ci-verbose"}, want: true, }, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { got := kgit.HasSkipCIPushOption(tt.pushOptions) if got != tt.want { t.Fatalf("hasSkipCIPushOption(%v) = %v, want %v", tt.pushOptions, got, tt.want) } }) } } func TestExecutorRoleBuildsMinimalSpindle(t *testing.T) { ctx := context.Background() dbPath := filepath.Join(t.TempDir(), "spindle.db") d, err := db.Make(ctx, dbPath) if err != nil { t.Fatalf("db.Make() error = %v", err) } cfg := &config.Config{Role: config.RoleExecutor} cfg.Server.DBPath = dbPath cfg.Server.Hostname = "executor.test" cfg.Server.Tap.Embed = true cfg.ArtifactStores.Disk.Dir = t.TempDir() cfg.Mill.ArtifactStore = "disk" s, err := New(ctx, cfg, d, map[string]models.Engine{}, nil) if err != nil { t.Fatalf("New() error = %v", err) } if s.jc != nil || s.tap != nil || s.e != nil || s.feed != nil || s.res != nil || s.vault != nil || s.scheduler != nil { t.Fatal("executor role built coordinator-only spindle dependencies") } rr := httptest.NewRecorder() req := httptest.NewRequest(http.MethodGet, "/", nil) s.Router().ServeHTTP(rr, req) if rr.Code != http.StatusOK { t.Fatalf("root status = %d, want %d", rr.Code, http.StatusOK) } rr = httptest.NewRecorder() req = httptest.NewRequest(http.MethodGet, "/_health", nil) s.Router().ServeHTTP(rr, req) if rr.Code != http.StatusOK { t.Fatalf("health status = %d, want %d", rr.Code, http.StatusOK) } if got := rr.Header().Get("Content-Type"); got != "application/json" { t.Fatalf("health content type = %q, want application/json", got) } if got := rr.Body.String(); got != "{\"status\":\"ok\"}\n" { t.Fatalf("health response = %q, want JSON ok status", got) } rr = httptest.NewRecorder() req = httptest.NewRequest(http.MethodGet, "/xrpc/_health", nil) s.Router().ServeHTTP(rr, req) if rr.Code != http.StatusNotFound { t.Fatalf("executor xrpc status = %d, want %d", rr.Code, http.StatusNotFound) } } type drainTestExecutor struct { connected chan struct{} connectCtx chan context.Context disconnected chan struct{} drainStarted chan struct{} releaseDrain chan struct{} } func (e *drainTestExecutor) Connect(ctx context.Context) { e.connectCtx <- ctx close(e.connected) <-ctx.Done() close(e.disconnected) } func (e *drainTestExecutor) Drain(ctx context.Context) error { close(e.drainStarted) select { case <-e.releaseDrain: return nil case <-ctx.Done(): return ctx.Err() } } func (e *drainTestExecutor) RegisterMetrics(*observability.Metrics) {} func (e *drainTestExecutor) SetQuotaClient(*executor.QuotaClient) {} func TestExecutorShutdownDrainsBeforeDisconnecting(t *testing.T) { exec := &drainTestExecutor{ connected: make(chan struct{}), connectCtx: make(chan context.Context, 1), disconnected: make(chan struct{}), drainStarted: make(chan struct{}), releaseDrain: make(chan struct{}, 1), } t.Cleanup(func() { select { case exec.releaseDrain <- struct{}{}: default: } }) cfg := &config.Config{Role: config.RoleExecutor} cfg.Server.ListenAddr = "127.0.0.1:0" cfg.Mill.DrainTimeout = 2 * time.Second s := &Spindle{ cfg: cfg, exec: exec, l: slog.New(slog.NewTextHandler(io.Discard, nil)), jobWake: make(chan struct{}, 1), } ctx, cancel := context.WithCancel(context.Background()) done := make(chan error, 1) go func() { done <- s.Start(ctx) }() select { case <-exec.connected: case <-time.After(2 * time.Second): t.Fatal("executor did not connect") } connectCtx := <-exec.connectCtx cancel() select { case <-exec.drainStarted: case <-time.After(2 * time.Second): t.Fatal("executor drain did not start") } if err := connectCtx.Err(); err != nil { t.Fatalf("Connect context canceled during drain: %v", err) } select { case <-exec.disconnected: t.Fatal("executor disconnected before drain completed") case err := <-done: t.Fatalf("Start returned before drain completed: %v", err) default: } exec.releaseDrain <- struct{}{} select { case <-exec.disconnected: case <-time.After(2 * time.Second): t.Fatal("executor did not disconnect after drain completed") } select { case err := <-done: if err != nil { t.Fatalf("Start() error = %v", err) } case <-time.After(2 * time.Second): t.Fatal("Start did not return after executor shutdown") } } func TestExecutorListenerFailureDrainsBeforeDisconnecting(t *testing.T) { exec := &drainTestExecutor{ connected: make(chan struct{}), connectCtx: make(chan context.Context, 1), disconnected: make(chan struct{}), drainStarted: make(chan struct{}), releaseDrain: make(chan struct{}, 1), } cfg := &config.Config{Role: config.RoleExecutor} cfg.Server.ListenAddr = "127.0.0.1:-1" cfg.Mill.DrainTimeout = 2 * time.Second s := &Spindle{ cfg: cfg, exec: exec, l: slog.New(slog.NewTextHandler(io.Discard, nil)), jobWake: make(chan struct{}, 1), } done := make(chan error, 1) go func() { done <- s.Start(context.Background()) }() <-exec.connected connectCtx := <-exec.connectCtx select { case <-exec.drainStarted: case <-time.After(2 * time.Second): t.Fatal("listener failure did not start executor drain") } if err := connectCtx.Err(); err != nil { t.Fatalf("Connect context canceled during drain: %v", err) } select { case <-exec.disconnected: t.Fatal("executor disconnected before drain completed") default: } exec.releaseDrain <- struct{}{} select { case err := <-done: if err == nil { t.Fatal("Start discarded the listener error") } case <-time.After(2 * time.Second): t.Fatal("Start did not return after listener failure cleanup") } } func TestNewOwnsLifecycleContextBeforeStart(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) dbPath := filepath.Join(t.TempDir(), "spindle.db") d, err := db.Make(ctx, dbPath) if err != nil { t.Fatal(err) } cfg := &config.Config{Role: config.RoleExecutor} cfg.Server.DBPath = dbPath cfg.Server.Hostname = "executor.test" cfg.ArtifactStores.Disk.Dir = t.TempDir() cfg.Mill.ArtifactStore = "disk" s, err := New(ctx, cfg, d, map[string]models.Engine{}, nil) if err != nil { t.Fatal(err) } t.Cleanup(s.rootCancel) cancel() if err := s.rootCtx.Err(); err != nil { t.Fatalf("lifecycle context inherited caller cancellation: %v", err) } }