package jobs import ( "context" "fmt" "log/slog" "os" "testing" "time" "git.nakama.town/fmartingr/hako/internal/archival/archiver" "git.nakama.town/fmartingr/hako/internal/extractors" "git.nakama.town/fmartingr/hako/internal/model" "github.com/stretchr/testify/require" ) func TestWorker_ProcessJob_Success(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() // Create mocks queue := NewMemoryQueue() mockArchiveService := &MockArchiveService{} mockArchiveStore := NewMockArchiveStore() mockArchiveFileStore := NewMockArchiveFileStore() mockLinkStore := NewMockLinkStore() logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelError})) // Create temporary storage for archiver manager tmpDir, err := os.MkdirTemp("", "test_worker_*") require.NoError(t, err) defer func() { _ = os.RemoveAll(tmpDir) }() archiverMgr := archiver.NewManager() extractorMgr := extractors.NewManager(logger) // Create worker worker := NewWorker(queue, archiverMgr, extractorMgr, mockArchiveService, mockArchiveStore, mockArchiveFileStore, mockLinkStore, tmpDir, logger) // Add a link to the mock store link := &model.Link{ ID: "link-123", URL: "https://example.com", UserID: "user-123", } mockLinkStore.AddLink(link) // Add an archive to the mock store archive := &model.Archive{ ID: "archive-123", LinkID: "link-123", UserID: "user-123", Status: model.ArchiveStatusPending, } mockArchiveStore.AddArchive(archive) // Enqueue a job payload := model.ArchiveLinkPayload{ LinkID: "link-123", ArchiveID: "archive-123", } err = queue.Enqueue(ctx, model.JobTypeArchiveLink, payload) require.NoError(t, err) // Start worker in background workerCtx, workerCancel := context.WithCancel(ctx) defer workerCancel() go func() { _ = worker.Start(workerCtx) }() // Wait for job to be processed time.Sleep(500 * time.Millisecond) // Verify the archive service was called require.Equal(t, 1, mockArchiveService.GetProcessedCount()) processedItems := mockArchiveService.GetProcessedItems() require.Len(t, processedItems, 1) require.Equal(t, "archive-123", processedItems[0]) } func TestWorker_ProcessJob_ArchiveServiceFailure(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() queue := NewMemoryQueue() mockArchiveService := NewMockArchiveService() mockArchiveService.ShouldFail = true mockArchiveStore := NewMockArchiveStore() mockLinkStore := NewMockLinkStore() logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelError})) tmpDir, err := os.MkdirTemp("", "test_worker_*") require.NoError(t, err) defer func() { _ = os.RemoveAll(tmpDir) }() extractorMgr := archiver.NewManager() contentExtractorMgr := extractors.NewManager(logger) mockArchiveFileStore := NewMockArchiveFileStore() worker := NewWorker(queue, extractorMgr, contentExtractorMgr, mockArchiveService, mockArchiveStore, mockArchiveFileStore, mockLinkStore, tmpDir, logger) link := &model.Link{ ID: "link-456", URL: "https://example.com/fail", UserID: "user-456", } mockLinkStore.AddLink(link) // Add archive to mock store archive := &model.Archive{ ID: "archive-456", LinkID: "link-456", UserID: "user-456", Status: model.ArchiveStatusPending, } mockArchiveStore.AddArchive(archive) payload := model.ArchiveLinkPayload{ LinkID: "link-456", ArchiveID: "archive-456", } err = queue.Enqueue(ctx, model.JobTypeArchiveLink, payload) require.NoError(t, err) workerCtx, workerCancel := context.WithCancel(ctx) defer workerCancel() go func() { _ = worker.Start(workerCtx) }() time.Sleep(500 * time.Millisecond) // Archive service should have been called but failed require.Equal(t, 1, mockArchiveService.GetProcessedCount()) } func TestWorker_ProcessJob_LinkNotFound(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() queue := NewMemoryQueue() mockArchiveService := NewMockArchiveService() mockArchiveStore := NewMockArchiveStore() mockLinkStore := NewMockLinkStore() logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelError})) tmpDir, err := os.MkdirTemp("", "test_worker_*") require.NoError(t, err) defer func() { _ = os.RemoveAll(tmpDir) }() extractorMgr := archiver.NewManager() contentExtractorMgr := extractors.NewManager(logger) mockArchiveFileStore := NewMockArchiveFileStore() worker := NewWorker(queue, extractorMgr, contentExtractorMgr, mockArchiveService, mockArchiveStore, mockArchiveFileStore, mockLinkStore, tmpDir, logger) // Don't add link to store - it won't be found payload := model.ArchiveLinkPayload{ LinkID: "nonexistent-link", ArchiveID: "archive-789", } err = queue.Enqueue(ctx, model.JobTypeArchiveLink, payload) require.NoError(t, err) workerCtx, workerCancel := context.WithCancel(ctx) defer workerCancel() go func() { _ = worker.Start(workerCtx) }() time.Sleep(500 * time.Millisecond) // Archive service should NOT have been called because link wasn't found require.Equal(t, 0, mockArchiveService.GetProcessedCount()) } func TestWorker_ProcessJob_InvalidPayload(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() queue := NewMemoryQueue() mockArchiveService := NewMockArchiveService() mockArchiveStore := NewMockArchiveStore() mockLinkStore := NewMockLinkStore() logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelError})) tmpDir, err := os.MkdirTemp("", "test_worker_*") require.NoError(t, err) defer func() { _ = os.RemoveAll(tmpDir) }() extractorMgr := archiver.NewManager() contentExtractorMgr := extractors.NewManager(logger) mockArchiveFileStore := NewMockArchiveFileStore() worker := NewWorker(queue, extractorMgr, contentExtractorMgr, mockArchiveService, mockArchiveStore, mockArchiveFileStore, mockLinkStore, tmpDir, logger) // Enqueue job with invalid payload (not valid JSON) err = queue.Enqueue(ctx, model.JobTypeArchiveLink, "not-valid-json{{{") require.NoError(t, err) workerCtx, workerCancel := context.WithCancel(ctx) defer workerCancel() go func() { _ = worker.Start(workerCtx) }() time.Sleep(500 * time.Millisecond) // Archive service should NOT have been called due to invalid payload require.Equal(t, 0, mockArchiveService.GetProcessedCount()) } func TestWorker_StartStop(t *testing.T) { ctx := context.Background() queue := NewMemoryQueue() mockArchiveService := NewMockArchiveService() mockArchiveStore := NewMockArchiveStore() mockLinkStore := NewMockLinkStore() logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelError})) tmpDir, err := os.MkdirTemp("", "test_worker_*") require.NoError(t, err) defer func() { _ = os.RemoveAll(tmpDir) }() extractorMgr := archiver.NewManager() contentExtractorMgr := extractors.NewManager(logger) mockArchiveFileStore := NewMockArchiveFileStore() worker := NewWorker(queue, extractorMgr, contentExtractorMgr, mockArchiveService, mockArchiveStore, mockArchiveFileStore, mockLinkStore, tmpDir, logger) // Start worker workerCtx, cancel := context.WithCancel(ctx) done := make(chan error, 1) go func() { done <- worker.Start(workerCtx) }() // Let it run briefly time.Sleep(100 * time.Millisecond) // Stop worker cancel() // Wait for worker to stop select { case err := <-done: require.NoError(t, err) case <-time.After(2 * time.Second): t.Fatal("Worker did not stop in time") } } func TestWorker_MultipleJobs(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() queue := NewMemoryQueue() mockArchiveService := NewMockArchiveService() mockArchiveStore := NewMockArchiveStore() mockLinkStore := NewMockLinkStore() logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelError})) tmpDir, err := os.MkdirTemp("", "test_worker_*") require.NoError(t, err) defer func() { _ = os.RemoveAll(tmpDir) }() extractorMgr := archiver.NewManager() contentExtractorMgr := extractors.NewManager(logger) mockArchiveFileStore := NewMockArchiveFileStore() worker := NewWorker(queue, extractorMgr, contentExtractorMgr, mockArchiveService, mockArchiveStore, mockArchiveFileStore, mockLinkStore, tmpDir, logger) // Add multiple links numJobs := 5 for i := 0; i < numJobs; i++ { linkID := fmt.Sprintf("link-%d", i) archiveID := fmt.Sprintf("archive-%d", i) link := &model.Link{ ID: linkID, URL: fmt.Sprintf("https://example.com/%d", i), UserID: "user-123", } mockLinkStore.AddLink(link) archive := &model.Archive{ ID: archiveID, LinkID: linkID, UserID: "user-123", Status: model.ArchiveStatusPending, } mockArchiveStore.AddArchive(archive) payload := model.ArchiveLinkPayload{ LinkID: linkID, ArchiveID: archiveID, } err = queue.Enqueue(ctx, model.JobTypeArchiveLink, payload) require.NoError(t, err) } // Start worker workerCtx, workerCancel := context.WithCancel(ctx) defer workerCancel() go func() { _ = worker.Start(workerCtx) }() // Wait for all jobs to be processed time.Sleep(1 * time.Second) // Verify all jobs were processed require.Equal(t, numJobs, mockArchiveService.GetProcessedCount()) processedItems := mockArchiveService.GetProcessedItems() require.Len(t, processedItems, numJobs) }