336 lines
9.3 KiB
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)
|
|
}
|