neoq/neoq_test.go

145 lines
2.6 KiB
Go
Raw Normal View History

2023-02-13 20:24:49 +00:00
package neoq
import (
"context"
"errors"
"fmt"
"log"
2023-02-19 23:52:37 +00:00
"os"
2023-02-13 20:24:49 +00:00
"testing"
"time"
)
2023-02-19 23:52:37 +00:00
var (
dbURL = "postgres://postgres:postgres@127.0.0.1:5432/neoq"
)
2023-02-13 20:24:49 +00:00
func TestWorkerListenConn(t *testing.T) {
const queue = "foobar"
2023-02-18 04:06:03 +00:00
ctx := context.TODO()
2023-02-19 23:52:37 +00:00
cnx := os.Getenv("DATABASE_URL")
if cnx == "" {
cnx = dbURL
}
pgBackend, err := NewPgBackend(ctx, cnx)
if err != nil {
t.Fatal(err)
}
2023-02-18 04:06:03 +00:00
nq, err := New(ctx, Backend(pgBackend))
2023-02-13 20:24:49 +00:00
if err != nil {
t.Fatal(err)
}
jobRan := false
numJobs := 1
var done = make(chan bool, numJobs)
handler := NewHandler(func(ctx context.Context) (err error) {
2023-02-16 23:15:39 +00:00
var j *Job
j, err = JobFromContext(ctx)
log.Println("queue:", j.Queue, "got job", "id:", j.ID, "messsage:", j.Payload["message"])
2023-02-13 20:24:49 +00:00
done <- true
return
})
handler = handler.
2023-02-18 18:45:32 +00:00
WithOption(HandlerDeadline(500 * time.Millisecond)).
WithOption(HandlerConcurreny(1))
2023-02-13 20:24:49 +00:00
if err != nil {
t.Error(err)
}
// Listen for jobs on the queue
2023-02-18 04:06:03 +00:00
nq.Listen(ctx, queue, handler)
2023-02-13 20:24:49 +00:00
// allow time for listener to start
time.Sleep(50 * time.Millisecond)
2023-02-13 20:24:49 +00:00
for i := 0; i < numJobs; i++ {
2023-02-18 04:06:03 +00:00
jid, err := nq.Enqueue(ctx, Job{
2023-02-13 20:24:49 +00:00
Queue: queue,
Payload: map[string]interface{}{
2023-02-16 23:15:39 +00:00
"message": fmt.Sprintf("hello world: %d", i),
2023-02-13 20:24:49 +00:00
},
})
if err != nil || jid == -1 {
t.Fatal("job was not enqueued. either it was duplicate or this error caused it:", err)
}
}
timeout := false
doneCnt := 0
for {
select {
case <-time.After(5 * time.Second):
timeout = true
err = errors.New("timed out waiting for job")
case <-done:
doneCnt++
}
if doneCnt >= numJobs {
jobRan = true
break
}
if timeout {
break
}
}
// Allow time for job status to be updated in the database
time.Sleep(50 * time.Millisecond)
if !jobRan {
t.Error(err)
}
}
2023-02-17 19:11:55 +00:00
func TestWorkerListenCron(t *testing.T) {
const cron = "* * * * * *"
2023-02-18 04:06:03 +00:00
ctx := context.TODO()
nq, err := New(ctx)
2023-02-17 19:11:55 +00:00
if err != nil {
t.Fatal(err)
}
jobRan := false
var done = make(chan bool)
handler := NewHandler(func(ctx context.Context) (err error) {
log.Println("got periodic job")
2023-02-17 19:11:55 +00:00
done <- true
return
})
2023-02-17 19:11:55 +00:00
handler = handler.
2023-02-18 18:45:32 +00:00
WithOption(HandlerDeadline(500 * time.Millisecond)).
WithOption(HandlerConcurreny(1))
2023-02-17 19:11:55 +00:00
if err != nil {
t.Error(err)
}
2023-02-18 04:06:03 +00:00
nq.ListenCron(ctx, cron, handler)
2023-02-17 19:11:55 +00:00
// allow time for listener to start
time.Sleep(50 * time.Millisecond)
select {
case <-time.After(5 * time.Second):
err = errors.New("timed out waiting for periodic job")
case <-done:
jobRan = true
}
// Allow time for job status to be updated in the database
time.Sleep(50 * time.Millisecond)
if !jobRan {
t.Error(err)
}
}