hako/internal/jobs/worker_test.go

336 lines
9.3 KiB
Go

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)
}