Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293package 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.ks != nil || s.res != nil || s.vault != 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) }}