mirror of
https://github.com/acaloiaro/neoq
synced 2026-07-21 18:29:08 +00:00
Move handlers.FromJobContext() -> jobs.FromContext() (#41)
This commit is contained in:
parent
d8c637852c
commit
81845a2303
10 changed files with 25 additions and 19 deletions
|
|
@ -47,7 +47,7 @@ Queue Handlers are simple Go functions that accept a `Context` parameter.
|
|||
ctx := context.Background()
|
||||
nq, _ := neoq.New(ctx)
|
||||
nq.Start(ctx, "hello_world", handler.New(func(ctx context.Context) (err error) {
|
||||
j, _ := handler.JobFromContext(ctx)
|
||||
j, _ := jobs.FromContext(ctx)
|
||||
log.Println("got job id:", j.ID, "messsage:", j.Payload["message"])
|
||||
return
|
||||
}))
|
||||
|
|
@ -82,7 +82,7 @@ nq, _ := neoq.New(ctx,
|
|||
)
|
||||
|
||||
nq.Start(ctx, "hello_world", handler.New(func(ctx context.Context) (err error) {
|
||||
j, _ := handler.JobFromContext(ctx)
|
||||
j, _ := jobs.FromContext(ctx)
|
||||
log.Println("got job id:", j.ID, "messsage:", j.Payload["message"])
|
||||
return
|
||||
}))
|
||||
|
|
|
|||
|
|
@ -518,6 +518,7 @@ func (p *PgBackend) moveToDeadQueue(ctx context.Context, tx pgx.Tx, j *jobs.Job,
|
|||
//
|
||||
// ultimately, this means that any time a database connection is lost while updating job status, then the job will be
|
||||
// processed at least one more time.
|
||||
// nolint: cyclop
|
||||
func (p *PgBackend) updateJob(ctx context.Context, jobErr error) (err error) {
|
||||
status := internal.JobStatusProcessed
|
||||
errMsg := ""
|
||||
|
|
@ -528,7 +529,7 @@ func (p *PgBackend) updateJob(ctx context.Context, jobErr error) (err error) {
|
|||
}
|
||||
|
||||
var job *jobs.Job
|
||||
if job, err = handler.JobFromContext(ctx); err != nil {
|
||||
if job, err = jobs.FromContext(ctx); err != nil {
|
||||
return fmt.Errorf("error getting job from context: %w", err)
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -21,8 +21,8 @@ const (
|
|||
type Config struct {
|
||||
BackendInitializer BackendInitializer
|
||||
ConnectionString string // a string containing connection details for the backend
|
||||
FutureJobWindow time.Duration // the window of time between now and RunAfter that goroutines are scheduled for future jobs
|
||||
JobCheckInterval time.Duration // the interval of time between checking for new future/retry jobs
|
||||
FutureJobWindow time.Duration // time duration between current time and job.RunAfter that goroutines schedule for future jobs
|
||||
IdleTransactionTimeout int // the number of milliseconds PgBackend transaction may idle before the connection is killed
|
||||
}
|
||||
|
||||
|
|
|
|||
3
env.sample
Normal file
3
env.sample
Normal file
|
|
@ -0,0 +1,3 @@
|
|||
TEST_DATABASE_URL=<TEST_DATABASE_URL>
|
||||
TEST_REDIS_URL=<TEST_REDIS_URL>
|
||||
|
||||
|
|
@ -27,7 +27,7 @@ func main() {
|
|||
// Option 1: add options when creating the handler
|
||||
h := handler.New(func(ctx context.Context) (err error) {
|
||||
var j *jobs.Job
|
||||
j, err = handler.JobFromContext(ctx)
|
||||
j, err = jobs.FromContext(ctx)
|
||||
log.Println("got job id:", j.ID, "messsage:", j.Payload["message"])
|
||||
done <- true
|
||||
return
|
||||
|
|
|
|||
|
|
@ -31,7 +31,7 @@ func main() {
|
|||
h := handler.New(func(ctx context.Context) (err error) {
|
||||
var j *jobs.Job
|
||||
time.Sleep(1 * time.Second)
|
||||
j, err = handler.JobFromContext(ctx)
|
||||
j, err = jobs.FromContext(ctx)
|
||||
log.Println("got job id:", j.ID, "messsage:", j.Payload["message"])
|
||||
done <- true
|
||||
return
|
||||
|
|
|
|||
|
|
@ -28,7 +28,7 @@ func main() {
|
|||
h := handler.New(func(ctx context.Context) (err error) {
|
||||
var j *jobs.Job
|
||||
time.Sleep(1 * time.Second)
|
||||
j, err = handler.JobFromContext(ctx)
|
||||
j, err = jobs.FromContext(ctx)
|
||||
log.Println("got job id:", j.ID, "messsage:", j.Payload["message"])
|
||||
done <- true
|
||||
return
|
||||
|
|
|
|||
|
|
@ -23,7 +23,7 @@ func main() {
|
|||
|
||||
h := handler.New(func(ctx context.Context) (err error) {
|
||||
var j *jobs.Job
|
||||
j, err = handler.JobFromContext(ctx)
|
||||
j, err = jobs.FromContext(ctx)
|
||||
log.Println("got job id:", j.ID, "messsage:", j.Payload["message"])
|
||||
return
|
||||
})
|
||||
|
|
|
|||
|
|
@ -7,7 +7,6 @@ import (
|
|||
"runtime"
|
||||
"time"
|
||||
|
||||
"github.com/acaloiaro/neoq/internal"
|
||||
"github.com/acaloiaro/neoq/jobs"
|
||||
)
|
||||
|
||||
|
|
@ -120,13 +119,3 @@ func Exec(ctx context.Context, handler Handler) (err error) {
|
|||
|
||||
return
|
||||
}
|
||||
|
||||
// JobFromContext fetches the job from a context if the job context variable is already set
|
||||
func JobFromContext(ctx context.Context) (j *jobs.Job, err error) {
|
||||
var ok bool
|
||||
if j, ok = ctx.Value(internal.JobCtxVarKey).(*jobs.Job); ok {
|
||||
return
|
||||
}
|
||||
|
||||
return nil, ErrContextHasNoJob
|
||||
}
|
||||
|
|
|
|||
13
jobs/jobs.go
13
jobs/jobs.go
|
|
@ -1,6 +1,7 @@
|
|||
package jobs
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/md5" // nolint: gosec
|
||||
"encoding/json"
|
||||
"errors"
|
||||
|
|
@ -8,10 +9,12 @@ import (
|
|||
"io"
|
||||
"time"
|
||||
|
||||
"github.com/acaloiaro/neoq/internal"
|
||||
"github.com/guregu/null"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrContextHasNoJob = errors.New("context has no Job")
|
||||
ErrJobTimeout = errors.New("timed out waiting for job(s)")
|
||||
ErrNoQueueSpecified = errors.New("this job does not specify a queue. please specify a queue")
|
||||
)
|
||||
|
|
@ -68,3 +71,13 @@ func FingerprintJob(j *Job) (err error) {
|
|||
|
||||
return
|
||||
}
|
||||
|
||||
// FromContext fetches the job from a context if the job context variable is set
|
||||
func FromContext(ctx context.Context) (j *Job, err error) {
|
||||
var ok bool
|
||||
if j, ok = ctx.Value(internal.JobCtxVarKey).(*Job); ok {
|
||||
return
|
||||
}
|
||||
|
||||
return nil, ErrContextHasNoJob
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue