package storage import ( "testing" "time" "github.com/darroyo/nasctl/internal/db" ) func TestJobManagerStartRejectsConcurrent(t *testing.T) { d, err := db.Open(":memory:") if err != nil { t.Fatalf("open db: %v", err) } defer d.Close() if err := d.Migrate(); err != nil { t.Fatalf("migrate: %v", err) } jm := NewJobManager(d, true) // Start a job that will block (no cmd specified — we just test rejection logic) _, err = jm.Start("snapraid_diff", nil) if err != nil { t.Fatalf("first Start: %v", err) } // Second start should be rejected _, err = jm.Start("snapraid_sync", nil) if err == nil { t.Error("second Start should have been rejected (job already running)") } } func TestJobManagerOk(t *testing.T) { d, _ := db.Open(":memory:") d.Migrate() defer d.Close() jmOn := NewJobManager(d, true) if !jmOn.Ok() { t.Error("Ok() = false, want true") } jmOff := NewJobManager(d, false) if jmOff.Ok() { t.Error("Ok() = true, want false") } } func TestJobManagerListEmpty(t *testing.T) { d, _ := db.Open(":memory:") d.Migrate() defer d.Close() jm := NewJobManager(d, false) jobs, err := jm.List(10, 0) if err != nil { t.Fatalf("List: %v", err) } if len(jobs) != 0 { t.Errorf("List = %d jobs, want 0", len(jobs)) } } func TestJobManagerResetOrphans(t *testing.T) { d, _ := db.Open(":memory:") d.Migrate() defer d.Close() jm := NewJobManager(d, false) // Create a running job directly in DB job, err := d.CreateStorageJob("snapraid_sync", "{}") if err != nil { t.Fatalf("CreateStorageJob: %v", err) } // Simulate orphaned running job _ = d.MarkJobRunning(job.ID, 12345) // Create a queued job _, _ = d.CreateStorageJob("snapraid_diff", "{}") if err := jm.ResetOrphans(); err != nil { t.Fatalf("ResetOrphans: %v", err) } jobs, _ := jm.List(10, 0) runningCount := 0 queuedCount := 0 for _, j := range jobs { if j.Status == "running" { runningCount++ } if j.Status == "queued" { queuedCount++ } if j.Status == "interrupted" { // should happen for the orphaned ones } } if runningCount != 0 || queuedCount != 0 { t.Errorf("after ResetOrphans: running=%d queued=%d, want both 0", runningCount, queuedCount) } } func TestJobManagerSubscribeAndBroadcast(t *testing.T) { d, _ := db.Open(":memory:") d.Migrate() defer d.Close() jm := NewJobManager(d, true) events, unsub := jm.Subscribe(9999) defer unsub() jm.broadcast(9999, Event{Type: "line", Line: "hello"}) jm.broadcast(9999, Event{Type: "end", Status: "success"}) select { case ev := <-events: if ev.Type != "line" || ev.Line != "hello" { t.Errorf("got event %+v, want {Type:line Line:hello}", ev) } case <-time.After(100 * time.Millisecond): t.Error("timeout waiting for event") } // second event select { case ev := <-events: if ev.Type != "end" || ev.Status != "success" { t.Errorf("got event %+v, want {Type:end Status:success}", ev) } case <-time.After(100 * time.Millisecond): t.Error("timeout waiting for end event") } }