diff --git a/cmd/server/tests/jobs_test.go b/cmd/server/tests/jobs_test.go new file mode 100644 index 0000000..33f1516 --- /dev/null +++ b/cmd/server/tests/jobs_test.go @@ -0,0 +1,232 @@ +package main + +import ( + "bytes" + "encoding/json" + "fmt" + "net/http" + "testing" + "time" + + "github.com/google/uuid" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestJobsHandler_CreateJob(t *testing.T) { + setup := setupTestServer(t) + defer setup.Close() + + token := setup.Token + + jobReq := map[string]interface{}{ + "type": "import", + "params": map[string]interface{}{ + "test": "data", + }, + } + + req, err := http.NewRequest("POST", setup.Server.URL+"/api/jobs", bytes.NewBuffer(jsonMarshal(jobReq))) + require.NoError(t, err) + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer "+token) + + client := &http.Client{} + resp, err := client.Do(req) + require.NoError(t, err) + defer resp.Body.Close() + + require.Equal(t, http.StatusAccepted, resp.StatusCode) + + var response map[string]interface{} + err = json.NewDecoder(resp.Body).Decode(&response) + require.NoError(t, err) + + assert.NotEmpty(t, response["job_id"]) + assert.Equal(t, "Job created", response["message"]) + assert.Equal(t, "import", response["type"]) + assert.Equal(t, "pending", response["status"]) +} + +func TestJobsHandler_CreateJob_InvalidType(t *testing.T) { + setup := setupTestServer(t) + defer setup.Close() + + token := setup.Token + + jobReq := map[string]interface{}{ + "type": "invalid_job_type", + "params": map[string]interface{}{}, + } + + req, err := http.NewRequest("POST", setup.Server.URL+"/api/jobs", bytes.NewBuffer(jsonMarshal(jobReq))) + require.NoError(t, err) + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer "+token) + + client := &http.Client{} + resp, err := client.Do(req) + require.NoError(t, err) + defer resp.Body.Close() + + require.Equal(t, http.StatusBadRequest, resp.StatusCode) +} + +func TestJobsHandler_GetJobStatus(t *testing.T) { + setup := setupTestServer(t) + defer setup.Close() + + token := setup.Token + + jobReq := map[string]interface{}{ + "type": "import", + "params": map[string]interface{}{}, + } + + createReq, err := http.NewRequest("POST", setup.Server.URL+"/api/jobs", bytes.NewBuffer(jsonMarshal(jobReq))) + require.NoError(t, err) + createReq.Header.Set("Content-Type", "application/json") + createReq.Header.Set("Authorization", "Bearer "+token) + + client := &http.Client{} + createResp, err := client.Do(createReq) + require.NoError(t, err) + defer createResp.Body.Close() + + var createResponse map[string]interface{} + err = json.NewDecoder(createResp.Body).Decode(&createResponse) + require.NoError(t, err) + + jobID := createResponse["job_id"].(string) + require.NotEmpty(t, jobID) + + getReq, err := http.NewRequest("GET", fmt.Sprintf("%s/api/jobs/%s", setup.Server.URL, jobID), nil) + require.NoError(t, err) + getReq.Header.Set("Authorization", "Bearer "+token) + + getResp, err := client.Do(getReq) + require.NoError(t, err) + defer getResp.Body.Close() + + require.Equal(t, http.StatusOK, getResp.StatusCode) + + var statusResponse map[string]interface{} + err = json.NewDecoder(getResp.Body).Decode(&statusResponse) + require.NoError(t, err) + + assert.Equal(t, jobID, statusResponse["id"]) + assert.Equal(t, "import", statusResponse["type"]) +} + +func TestJobsHandler_GetJobStatus_NotFound(t *testing.T) { + setup := setupTestServer(t) + defer setup.Close() + + token := setup.Token + + fakeJobID := uuid.New().String() + + req, err := http.NewRequest("GET", fmt.Sprintf("%s/api/jobs/%s", setup.Server.URL, fakeJobID), nil) + require.NoError(t, err) + req.Header.Set("Authorization", "Bearer "+token) + + client := &http.Client{} + resp, err := client.Do(req) + require.NoError(t, err) + defer resp.Body.Close() + + require.Equal(t, http.StatusNotFound, resp.StatusCode) +} + +func TestJobsHandler_CreateAndTrackJob(t *testing.T) { + setup := setupTestServer(t) + defer setup.Close() + + token := setup.Token + + jobReq := map[string]interface{}{ + "type": "import", + "params": map[string]interface{}{}, + } + + createReq, err := http.NewRequest("POST", setup.Server.URL+"/api/jobs", bytes.NewBuffer(jsonMarshal(jobReq))) + require.NoError(t, err) + createReq.Header.Set("Content-Type", "application/json") + createReq.Header.Set("Authorization", "Bearer "+token) + + client := &http.Client{} + createResp, err := client.Do(createReq) + require.NoError(t, err) + defer createResp.Body.Close() + + var createResponse map[string]interface{} + err = json.NewDecoder(createResp.Body).Decode(&createResponse) + require.NoError(t, err) + + jobID := createResponse["job_id"].(string) + + var finalStatus string + for i := 0; i < 20; i++ { + time.Sleep(500 * time.Millisecond) + + getReq, err := http.NewRequest("GET", fmt.Sprintf("%s/api/jobs/%s", setup.Server.URL, jobID), nil) + require.NoError(t, err) + getReq.Header.Set("Authorization", "Bearer "+token) + + getResp, err := client.Do(getReq) + require.NoError(t, err) + + var statusResponse map[string]interface{} + err = json.NewDecoder(getResp.Body).Decode(&statusResponse) + getResp.Body.Close() + require.NoError(t, err) + + if statusResponse["status"] != nil { + finalStatus = statusResponse["status"].(string) + if finalStatus == "completed" || finalStatus == "failed" { + break + } + } + } + + assert.True(t, finalStatus == "completed" || finalStatus == "failed", + fmt.Sprintf("Job should complete, got status: %s", finalStatus)) +} + +func TestJobsHandler_CreateJob_Unauthorized(t *testing.T) { + setup := setupTestServer(t) + defer setup.Close() + + jobReq := map[string]interface{}{ + "type": "import", + "params": map[string]interface{}{}, + } + + req, err := http.NewRequest("POST", setup.Server.URL+"/api/jobs", bytes.NewBuffer(jsonMarshal(jobReq))) + require.NoError(t, err) + req.Header.Set("Content-Type", "application/json") + + client := &http.Client{} + resp, err := client.Do(req) + require.NoError(t, err) + defer resp.Body.Close() + + require.Equal(t, http.StatusUnauthorized, resp.StatusCode) +} + +func TestJobsHandler_GetJobStatus_Unauthorized(t *testing.T) { + setup := setupTestServer(t) + defer setup.Close() + + jobID := uuid.New().String() + + req, err := http.NewRequest("GET", fmt.Sprintf("%s/api/jobs/%s", setup.Server.URL, jobID), nil) + require.NoError(t, err) + + client := &http.Client{} + resp, err := client.Do(req) + require.NoError(t, err) + defer resp.Body.Close() + + require.Equal(t, http.StatusUnauthorized, resp.StatusCode) +} diff --git a/cmd/server/tests/worker_test.go b/cmd/server/tests/worker_test.go new file mode 100644 index 0000000..a9e81c3 --- /dev/null +++ b/cmd/server/tests/worker_test.go @@ -0,0 +1,293 @@ +package main + +import ( + "bookhoard/internal/database" + "bookhoard/internal/services" + "bytes" + "context" + "encoding/json" + "fmt" + "net/http" + "os" + "path/filepath" + "testing" + "time" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgtype" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestWorker_DirectoryScanJob(t *testing.T) { + setup := setupTestServer(t) + defer setup.Close() + + token := setup.Token + + // Create library + createLibReq := map[string]interface{}{ + "name": "Worker Test Library", + "description": "Test library for worker", + "type": "ebooks", + } + createLibURL := setup.Server.URL + "/api/libraries" + createLibReqHTTP, _ := http.NewRequest("POST", createLibURL, bytes.NewBuffer(jsonMarshal(createLibReq))) + createLibReqHTTP.Header.Set("Content-Type", "application/json") + createLibReqHTTP.Header.Set("Authorization", "Bearer "+token) + + client := &http.Client{} + createLibResp, err := client.Do(createLibReqHTTP) + require.NoError(t, err) + require.Equal(t, http.StatusCreated, createLibResp.StatusCode) + + var createLibResponse map[string]interface{} + json.NewDecoder(createLibResp.Body).Decode(&createLibResponse) + createLibResp.Body.Close() + + libraryID, ok := createLibResponse["id"].(string) + require.True(t, ok) + require.NotEmpty(t, libraryID) + + // Create a temporary directory for testing + tmpDir := t.TempDir() + + // Add folder to library + folderURL := fmt.Sprintf("%s/api/libraries/%s/folders", setup.Server.URL, libraryID) + folderReq := map[string]interface{}{ + "folder_path": tmpDir, + } + req, _ := http.NewRequest("POST", folderURL, bytes.NewBuffer(jsonMarshal(folderReq))) + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer "+token) + + resp, err := client.Do(req) + require.NoError(t, err) + resp.Body.Close() + require.Equal(t, http.StatusCreated, resp.StatusCode) + + // Create test files in the directory + for i := 0; i < 3; i++ { + filePath := filepath.Join(tmpDir, fmt.Sprintf("book%d.epub", i)) + err := os.WriteFile(filePath, []byte(fmt.Sprintf("test content %d", i)), 0644) + require.NoError(t, err) + } + + // Submit directory scan job through WorkerInstance + job := &services.Job{ + ID: uuid.New().String(), + Type: services.JobTypeDirectoryScan, + Params: map[string]interface{}{ + "directory": tmpDir, + "db": setup.DB, + }, + Status: services.JobStatusPending, + } + + services.WorkerInstance.Enqueue(job) + + // Wait for job to process + time.Sleep(2 * time.Second) + + // Check that items were created in database + ctx := context.Background() + items, err := setup.DB.ListMediaItems(ctx, database.ListMediaItemsParams{}) + require.NoError(t, err) + + // Should have at least 3 items (our test files) + greaterOrEqual := func(a, b int) bool { + return a >= b + } + assert.True(t, greaterOrEqual(len(items), 3), fmt.Sprintf("Expected at least 3 items, got %d", len(items))) +} + +func TestWorker_SetFoldersJob(t *testing.T) { + setup := setupTestServer(t) + defer setup.Close() + + token := setup.Token + + // Create library + createLibReq := map[string]interface{}{ + "name": "Set Folders Test Library", + "description": "Test library for set folders", + "type": "ebooks", + } + createLibURL := setup.Server.URL + "/api/libraries" + createLibReqHTTP, _ := http.NewRequest("POST", createLibURL, bytes.NewBuffer(jsonMarshal(createLibReq))) + createLibReqHTTP.Header.Set("Content-Type", "application/json") + createLibReqHTTP.Header.Set("Authorization", "Bearer "+token) + + client := &http.Client{} + createLibResp, err := client.Do(createLibReqHTTP) + require.NoError(t, err) + require.Equal(t, http.StatusCreated, createLibResp.StatusCode) + + var createLibResponse map[string]interface{} + json.NewDecoder(createLibResp.Body).Decode(&createLibResponse) + createLibResp.Body.Close() + + libraryID, ok := createLibResponse["id"].(string) + require.True(t, ok) + require.NotEmpty(t, libraryID) + + // Create temporary directory + tmpDir := t.TempDir() + + // Submit set folders job through WorkerInstance + job := &services.Job{ + ID: uuid.New().String(), + Type: services.JobTypeSetFolders, + Params: map[string]interface{}{ + "folders": []string{tmpDir}, + "db": setup.DB, + }, + Status: services.JobStatusPending, + } + + services.WorkerInstance.Enqueue(job) + + // Wait for job to process + time.Sleep(1 * time.Second) + + // Verify folder was added to library + ctx := context.Background() + folders, err := setup.DB.GetLibraryFolders(ctx, pgtype.UUID{Bytes: [16]byte(uuid.MustParse(libraryID)), Valid: true}) + require.NoError(t, err) + assert.NotEmpty(t, folders, "Library should have folders") + + // Find our folder + found := false + for _, folder := range folders { + if folder.FolderPath == tmpDir { + found = true + break + } + } + assert.True(t, found, "Folder should be in library") +} + +func TestWorker_ConcurrentJobs(t *testing.T) { + setup := setupTestServer(t) + defer setup.Close() + + ctx := context.Background() + + // Create a library and get admin user + adminID := getTestUserID(t, setup.DB) + + libraryType, err := setup.DB.GetLibraryTypeByName(ctx, "ebooks") + require.NoError(t, err) + + library, err := setup.DB.CreateLibrary(ctx, database.CreateLibraryParams{ + Name: "Concurrent Job Test Library", + Description: pgtype.Text{String: "Test library", Valid: true}, + LibraryTypeID: libraryType.ID, + CreatedByAdminID: pgtype.UUID{Bytes: [16]byte(adminID), Valid: true}, + }) + require.NoError(t, err) + + // Create multiple temporary directories + var tmpDirs []string + for i := 0; i < 3; i++ { + tmpDir := t.TempDir() + tmpDirs = append(tmpDirs, tmpDir) + + // Add folder to library + _, err = setup.DB.AddLibraryFolder(ctx, database.AddLibraryFolderParams{ + LibraryID: library.ID, + FolderPath: tmpDir, + }) + require.NoError(t, err) + } + + // Submit multiple concurrent jobs + for _, tmpDir := range tmpDirs { + job := &services.Job{ + ID: uuid.New().String(), + Type: services.JobTypeDirectoryScan, + Params: map[string]interface{}{ + "directory": tmpDir, + "db": setup.DB, + }, + Status: services.JobStatusPending, + } + services.WorkerInstance.Enqueue(job) + } + + // Wait for all jobs to process + time.Sleep(3 * time.Second) + + // Verify all directories were processed + allItems, err := setup.DB.ListMediaItems(ctx, database.ListMediaItemsParams{}) + require.NoError(t, err) + assert.NotEmpty(t, allItems, "Should have media items from concurrent jobs") +} + +func TestWorker_JobStatusTracking(t *testing.T) { + setup := setupTestServer(t) + defer setup.Close() + + ctx := context.Background() + + // Create test library + adminID := getTestUserID(t, setup.DB) + libraryType, err := setup.DB.GetLibraryTypeByName(ctx, "ebooks") + require.NoError(t, err) + + library, err := setup.DB.CreateLibrary(ctx, database.CreateLibraryParams{ + Name: "Job Status Test Library", + Description: pgtype.Text{String: "Test library", Valid: true}, + LibraryTypeID: libraryType.ID, + CreatedByAdminID: pgtype.UUID{Bytes: [16]byte(adminID), Valid: true}, + }) + require.NoError(t, err) + + tmpDir := t.TempDir() + _, err = setup.DB.AddLibraryFolder(ctx, database.AddLibraryFolderParams{ + LibraryID: library.ID, + FolderPath: tmpDir, + }) + require.NoError(t, err) + + // Submit job and track status + jobID := uuid.New().String() + job := &services.Job{ + ID: jobID, + Type: services.JobTypeDirectoryScan, + Params: map[string]interface{}{ + "directory": tmpDir, + "db": setup.DB, + }, + Status: services.JobStatusPending, + } + + services.WorkerInstance.Enqueue(job) + + // Poll job status + var finalStatus string + for i := 0; i < 10; i++ { + time.Sleep(500 * time.Millisecond) + + result, exists := services.WorkerInstance.GetJobStatus(jobID) + if exists { + finalStatus = string(result.Status) + if result.Status == services.JobStatusCompleted || result.Status == services.JobStatusFailed { + break + } + } + } + + // Job should complete (success or failure) + assert.True(t, finalStatus == "completed" || finalStatus == "failed", + fmt.Sprintf("Job should complete, got status: %s", finalStatus)) +} + +// Helper function to marshal JSON +func jsonMarshal(v interface{}) []byte { + b, err := json.Marshal(v) + if err != nil { + panic(err) + } + return b +}