package web import ( "context" "errors" "fmt" "gitgud.io/mike/mpv-manager/pkg/constants" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "strings" "sync" "sync/atomic" "testing" "time" "go.uber.org/goleak" ) // TestJobSnapshot_ConcurrentWriters exercises JobSnapshot and the list // accessors under concurrent mutation. Run with -race to detect data races. func TestJobSnapshot_ConcurrentWriters(t *testing.T) { defer goleak.VerifyNone(t) jm := NewJobManager() job := jm.CreateJob("install", "mpv-test-snapshot", "MPV") var wg sync.WaitGroup stop := make(chan struct{}) // Concurrent writers mutating the job for i := 0; i < 4; i++ { wg.Add(1) go func(n int) { defer wg.Done() for { select { case <-stop: return default: jm.AddOutput(job.ID, fmt.Sprintf("line from writer %d", n)) jm.UpdateProgress(job.ID, n, "working") } } }(i) } // Concurrent readers taking snapshots for i := 0; i < 4; i++ { wg.Add(1) go func() { defer wg.Done() for { select { case <-stop: return default: snap, ok := jm.JobSnapshot(job.ID) if !ok { t.Error("JobSnapshot: job not found") return } _ = snap.Status _ = len(snap.Output) _ = jobToAPIResponse(snap) _ = jm.GetActiveJobs() _ = jm.GetRecentJobs() _, _ = jm.jobSummaries() } } }() } time.Sleep(100 * time.Millisecond) close(stop) wg.Wait() } // TestJobSnapshot_ReturnsCopy verifies that mutating a snapshot does not // affect the job stored in the manager (copied struct and Output slice). func TestJobSnapshot_ReturnsCopy(t *testing.T) { jm := NewJobManager() job := jm.CreateJob("install", "mpv-test-copy", "MPV") jm.AddOutput(job.ID, "line one") snap, ok := jm.JobSnapshot(job.ID) if !ok { t.Fatal("JobSnapshot: job not found") } if len(snap.Output) != 1 { t.Fatalf("Expected 1 output line, got %d", len(snap.Output)) } // Mutate the snapshot snap.Output[0] = "corrupted" snap.Status = "bogus" snap2, ok := jm.JobSnapshot(job.ID) if !ok { t.Fatal("JobSnapshot: job not found after mutation") } if snap2.Output[0] == "corrupted" { t.Error("Snapshot shares the Output backing array with the stored job") } if snap2.Status == "bogus" { t.Error("Snapshot shares state with the stored job") } } // TestJobSnapshot_NotFound verifies the missing-job case. func TestJobSnapshot_NotFound(t *testing.T) { jm := NewJobManager() if _, ok := jm.JobSnapshot("nonexistent"); ok { t.Error("Expected JobSnapshot to return false for unknown job ID") } } // TestCreateJobIfNoneActive_Conflict verifies the atomic check-and-create // rejects duplicates for the same method while a job is active. func TestCreateJobIfNoneActive_Conflict(t *testing.T) { defer goleak.VerifyNone(t) jm := NewJobManager() job, created := jm.CreateJobIfNoneActive("install", "mpv-test", "MPV") if !created || job == nil { t.Fatal("Expected first create to succeed") } if !jm.HasActiveJob("mpv-test") { t.Error("Expected HasActiveJob to report the active job") } // Duplicate while active -> conflict if dup, ok := jm.CreateJobIfNoneActive("install", "mpv-test", "MPV"); ok || dup != nil { t.Error("Expected duplicate create to be rejected while job is active") } // Different method is unaffected if _, ok := jm.CreateJobIfNoneActive("install", "mpv-test-other", "MPV Other"); !ok { t.Error("Expected create for a different method to succeed") } // After completion, creating for the same method succeeds again jm.CompleteJob(job.ID) if jm.HasActiveJob("mpv-test") { t.Error("Expected HasActiveJob to be false after completion") } if _, ok := jm.CreateJobIfNoneActive("install", "mpv-test", "MPV"); !ok { t.Error("Expected create to succeed after the previous job completed") } } // TestCreateJobIfNoneActive_Concurrent verifies that racing creators cannot // produce duplicate active jobs for the same method (TOCTOU fix). func TestCreateJobIfNoneActive_Concurrent(t *testing.T) { defer goleak.VerifyNone(t) jm := NewJobManager() const creators = 16 var wg sync.WaitGroup var successes int32 for i := 0; i < creators; i++ { wg.Add(1) go func() { defer wg.Done() if _, ok := jm.CreateJobIfNoneActive("install", "mpv-race", "MPV"); ok { atomic.AddInt32(&successes, 1) } }() } wg.Wait() if successes != 1 { t.Errorf("Expected exactly 1 successful create under concurrency, got %d", successes) } } func TestCancelJobWaitsForWorkerAndRetainsResources(t *testing.T) { jm := NewJobManager() job, created := jm.CreateJobWithResources("install", "method-a", "MPV", "config:shared", "target:a") if !created { t.Fatal("expected first job to acquire resources") } if err := jm.CancelJob(job.ID); err != nil { t.Fatalf("CancelJob() error = %v", err) } snapshot, ok := jm.JobSnapshot(job.ID) if !ok || snapshot.Status != "cancelling" || snapshot.CompletedAt != nil { t.Fatalf("cancellation request archived or completed the job: %+v, exists=%v", snapshot, ok) } if len(jm.GetRecentJobs()) != 0 { t.Fatal("cancelling job must not enter history before worker acknowledgement") } if conflicting, ok := jm.CreateJobWithResources("install", "method-b", "Other", "config:shared"); ok || conflicting != nil { t.Fatal("cancelling worker released its shared resource prematurely") } jm.SetError(job.ID, context.Canceled, "worker stopped at a safe boundary") recent := jm.GetRecentJobs() if len(recent) != 1 || recent[0].Status != "cancelled" { t.Fatalf("worker acknowledgement did not archive cancellation: %+v", recent) } if replacement, ok := jm.CreateJobWithResources("install", "method-b", "Other", "config:shared"); !ok || replacement == nil { t.Fatal("terminal worker acknowledgement did not release shared resource") } } func TestCancelJobRejectedAfterCommitBoundary(t *testing.T) { jm := NewJobManager() job, ok := jm.CreateJobWithResources("uninstall", "method-a", "MPV", "target:a") if !ok { t.Fatal("create job failed") } if err := jm.BeginCommit(job.ID); err != nil { t.Fatalf("BeginCommit() error = %v", err) } if err := jm.CancelJob(job.ID); !errors.Is(err, ErrJobCommitStarted) { t.Fatalf("CancelJob() error = %v, want ErrJobCommitStarted", err) } if job.Context().Err() != nil { t.Fatalf("rejected cancellation changed worker context: %v", job.Context().Err()) } jm.CompleteJob(job.ID) recent := jm.GetRecentJobs() if len(recent) != 1 || recent[0].Status != "complete" { t.Fatalf("commit-bound job did not complete authoritatively: %+v", recent) } } func TestPartialSuccessIsTerminalAndReleasesResources(t *testing.T) { jm := NewJobManager() job, ok := jm.CreateJobWithResources("install", "method-a", "MPV", "config:manager") if !ok { t.Fatal("create job failed") } jm.SetPartialSuccess(job.ID, errors.New("disk full"), "installed but untracked") recent := jm.GetRecentJobs() if len(recent) != 1 || recent[0].Status != "partial" || !recent[0].IsComplete() || !recent[0].IsError() { t.Fatalf("partial result is not a terminal non-success: %+v", recent) } if next, ok := jm.CreateJobWithResources("install", "method-b", "Other", "config:manager"); !ok || next == nil { t.Fatal("partial terminal result did not release resources") } } func TestTryAcquireResourcesIsAtomicAndIdempotent(t *testing.T) { jm := NewJobManager() release, ok := jm.TryAcquireResources("request-a", "b", "a", "a") if !ok { t.Fatal("first request should acquire deduplicated resources") } if _, ok := jm.TryAcquireResources("request-b", "c", "b"); ok { t.Fatal("multi-resource acquisition must fail atomically on one conflict") } if releaseC, ok := jm.TryAcquireResources("request-c", "c"); !ok { t.Fatal("failed acquisition leaked its otherwise-free resource") } else { releaseC() } release() release() if releaseB, ok := jm.TryAcquireResources("request-b", "a", "b"); !ok { t.Fatal("idempotent release did not free resources") } else { releaseB() } } // TestBroadcastEvent_TerminalEventDeliveredWhenBufferFull verifies that // terminal events (complete/status) are never dropped: when a client's buffer // is full, the oldest events are drained to make room. func TestBroadcastEvent_TerminalEventDeliveredWhenBufferFull(t *testing.T) { defer goleak.VerifyNone(t) jm := NewJobManager() client := jm.RegisterClient() defer jm.UnregisterClient(client) // Fill the client buffer with droppable output events for i := 0; i < cap(client); i++ { jm.broadcastEvent(JobEvent{JobID: "job1", Type: "output", Data: JobOutputData{Line: "noise"}}) } // Terminal event must displace the oldest buffered event, not be dropped jm.broadcastEvent(JobEvent{JobID: "job1", Type: "complete", Data: JobCompleteData{Success: true}}) drained := 0 sawComplete := false for { select { case ev := <-client: drained++ if ev.Type == "complete" { sawComplete = true } default: if !sawComplete { t.Error("Terminal 'complete' event was dropped from a full client buffer") } if drained != cap(client) { t.Errorf("Expected buffer to stay at capacity %d, drained %d events", cap(client), drained) } return } } } // TestBroadcastEvent_StatusEventDeliveredWhenBufferFull verifies "status" // events (error/cancelled) also survive a full buffer. func TestBroadcastEvent_StatusEventDeliveredWhenBufferFull(t *testing.T) { defer goleak.VerifyNone(t) jm := NewJobManager() client := jm.RegisterClient() defer jm.UnregisterClient(client) for i := 0; i < cap(client); i++ { jm.broadcastEvent(JobEvent{JobID: "job1", Type: "progress", Data: JobProgressData{Progress: 1}}) } jm.broadcastEvent(JobEvent{JobID: "job1", Type: "status", Data: JobStatusData{Status: "error"}}) sawStatus := false for { select { case ev := <-client: if ev.Type == "status" { sawStatus = true } default: if !sawStatus { t.Error("Terminal 'status' event was dropped from a full client buffer") } return } } } // TestBroadcastEvent_OutputEventDroppedWhenBufferFull verifies that // non-terminal events are still droppable (slow clients must not block). func TestBroadcastEvent_OutputEventDroppedWhenBufferFull(t *testing.T) { defer goleak.VerifyNone(t) jm := NewJobManager() client := jm.RegisterClient() defer jm.UnregisterClient(client) // Fill the buffer, then overflow it with droppable events for i := 0; i < cap(client)+10; i++ { jm.broadcastEvent(JobEvent{JobID: "job1", Type: "output", Data: JobOutputData{Line: "noise"}}) } drained := 0 for { select { case <-client: drained++ default: if drained != cap(client) { t.Errorf("Expected droppable events to be skipped at capacity %d, drained %d", cap(client), drained) } return } } } func TestCloseClientsIsIdempotentAndRejectsNewStreams(t *testing.T) { jm := NewJobManager() client := jm.RegisterClient() jm.CloseClients() jm.CloseClients() if _, ok := <-client; ok { t.Fatal("registered client was not closed") } lateClient := jm.RegisterClient() if _, ok := <-lateClient; ok { t.Fatal("client registered after broadcaster shutdown remained open") } // A worker racing shutdown must be able to publish without panicking. jm.broadcastEvent(JobEvent{JobID: "late", Type: "complete", Data: JobCompleteData{Success: true}}) jm.UnregisterClient(client) jm.UnregisterClient(lateClient) } func TestRegisterClientEnforcesConnectionLimit(t *testing.T) { jm := NewJobManager() clients := make([]chan JobEvent, 0, maxSSEClients) for range maxSSEClients { client := jm.RegisterClient() if client == nil { t.Fatal("client limit was enforced before capacity") } clients = append(clients, client) } if client := jm.RegisterClient(); client != nil { t.Fatal("client beyond capacity was accepted") } jm.UnregisterClient(clients[0]) replacement := jm.RegisterClient() if replacement == nil { t.Fatal("released client capacity was not reusable") } jm.UnregisterClient(replacement) for _, client := range clients[1:] { jm.UnregisterClient(client) } } // TestJobIsComplete verifies terminal-state detection including cancelled. func TestJobIsComplete(t *testing.T) { tests := []struct { status string expected bool }{ {"pending", false}, {"running", false}, {"complete", true}, {"error", true}, {"partial", true}, {"cancelled", true}, {"cancelling", false}, } for _, tt := range tests { t.Run(tt.status, func(t *testing.T) { job := &Job{Status: tt.status} if got := job.IsComplete(); got != tt.expected { t.Errorf("IsComplete() with status %q = %v, want %v", tt.status, got, tt.expected) } }) } } // TestAddOutputCapsStoredLines verifies that Job.Output is bounded to // maxJobOutputLines, dropping the oldest lines. func TestAddOutputCapsStoredLines(t *testing.T) { defer goleak.VerifyNone(t) jm := NewJobManager() job := jm.CreateJob("install", "mpv-test-cap", "MPV") total := maxJobOutputLines + 500 for i := 0; i < total; i++ { jm.AddOutput(job.ID, fmt.Sprintf("line %d", i)) } snap, ok := jm.JobSnapshot(job.ID) if !ok { t.Fatal("JobSnapshot: job not found") } if len(snap.Output) != maxJobOutputLines { t.Fatalf("Expected output capped at %d lines, got %d", maxJobOutputLines, len(snap.Output)) } // One retained line is reserved for the visible truncation notice. if !strings.Contains(snap.Output[0], "Earlier output omitted") || !strings.Contains(snap.Output[1], "line 501") { t.Errorf("Expected a truncation notice and newest tail, first lines: %q", snap.Output[:2]) } // Newest line must still be present last := snap.Output[len(snap.Output)-1] if !strings.Contains(last, fmt.Sprintf("line %d", total-1)) { t.Errorf("Expected newest line to be kept, last stored line: %q", last) } } func TestAddOutputBoundsBytesAndSSELines(t *testing.T) { jm := NewJobManager() job := jm.CreateJob("install", "fixture", "fixture") for index := 0; index < 300; index++ { jm.AddOutput(job.ID, strings.Repeat("x", constants.MaxOutputLineBytes*2)) } jm.AddOutput(job.ID, "last output") snapshot, ok := jm.JobSnapshot(job.ID) require.True(t, ok) storedBytes := 0 for _, line := range snapshot.Output { storedBytes += len(line) assert.LessOrEqual(t, len(line), constants.MaxOutputLineBytes+11, "timestamped SSE output must remain bounded") } assert.LessOrEqual(t, storedBytes, constants.MaxJobOutputBytes) assert.Equal(t, constants.OutputTailMarker, snapshot.Output[0]) assert.Contains(t, snapshot.Output[len(snapshot.Output)-1], "last output") assert.Less(t, len(snapshot.Output), 200, "the byte limit must take effect before the line limit") } func TestJobProgressReservesCompletionUntilWorkerFinishes(t *testing.T) { jm := NewJobManager() job := jm.CreateJob("install", "mpv-binary", "MPV") jm.AddOutput(job.ID, "Downloading archive: 100%") snapshot, ok := jm.JobSnapshot(job.ID) require.True(t, ok) require.Equal(t, 99, snapshot.Progress) require.False(t, snapshot.IsComplete()) jm.UpdateProgress(job.ID, 100, "Extracting archive (100%)") snapshot, _ = jm.JobSnapshot(job.ID) require.Equal(t, 99, snapshot.Progress) jm.CompleteJob(job.ID) jm.UpdateProgress(job.ID, 10, "late progress") jm.AddOutput(job.ID, "late output: 100%") recent := jm.GetRecentJobs() require.Len(t, recent, 1) snapshot = recent[0] require.Equal(t, "complete", snapshot.Status) require.Equal(t, 100, snapshot.Progress) } func TestJobSummariesRetainMetadataAndDoNotExposeOutput(t *testing.T) { jm := NewJobManager() job := jm.CreateJob("adopt", "fixture", "Summary fixture") jm.AddOutput(job.ID, "retained details output") jm.UpdateProgress(job.ID, 42, "Preparing") active, recent := jm.jobSummaries() require.Len(t, active, 1) require.Empty(t, recent) assert.Nil(t, active[0].Output) assert.Equal(t, "adopt", active[0].Type) assert.Equal(t, 42, active[0].Progress) jm.CompleteJob(job.ID) active, recent = jm.jobSummaries() require.Empty(t, active) require.Len(t, recent, 1) assert.Nil(t, recent[0].Output) assert.Equal(t, "complete", recent[0].Status) assert.NotNil(t, recent[0].CompletedAt) assert.Contains(t, jm.GetRecentJobs()[0].Output[0], "retained details output") } func BenchmarkJobListSnapshots(b *testing.B) { output := make([]string, 2000) for i := range output { output[i] = strings.Repeat("x", 500) } jm := &JobManager{jobs: map[string]*Job{"active": {ID: "active", Status: "running", Output: output}}} for i := 0; i < 10; i++ { jm.recentJobs = append(jm.recentJobs, &Job{Status: "complete", Output: output}) } b.Run("metadata-only", func(b *testing.B) { b.ReportAllocs() for b.Loop() { _, _ = jm.jobSummaries() } }) b.Run("output-copies", func(b *testing.B) { b.ReportAllocs() for b.Loop() { _ = jm.GetActiveJobs() _ = jm.GetRecentJobs() } }) }