Add tests for Jobs API and Worker job processing
- Add jobs_test.go with tests for job creation and status retrieval - Add worker_test.go with tests for job processing
This commit is contained in:
@@ -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)
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
Reference in New Issue
Block a user