This commit is contained in:
Pranav Shridhar 2024-04-09 17:14:24 +05:30
parent 330893cc06
commit ddbc6b1cbf

View file

@ -6,6 +6,7 @@ import (
"embed"
"errors"
"fmt"
"log"
"os"
"strings"
"sync"
@ -303,6 +304,8 @@ func (s *SqliteBackend) getAllPendingJobIDs(ctx context.Context, queue string) (
}
func (s *SqliteBackend) handleJob(ctx context.Context, jobID string) (err error) {
log.Println("PS::handleJob start")
job, err := s.getJobWithoutTx(ctx, jobID)
if err != nil {
s.logger.Error("could not retrieve row with job id", slog.String("err", err.Error()))
@ -323,9 +326,9 @@ func (s *SqliteBackend) handleJob(ctx context.Context, jobID string) (err error)
return handler.ErrNoHandlerForQueue
}
fmt.Printf("Current job : %+v", job)
log.Printf("Current job : %+v", job)
ctx = withJobContext(ctx, job)
fmt.Printf("Updated ctx with job : %+v", ctx)
log.Printf("Updated ctx with job : %+v", ctx)
var tx *sql.Tx
s.dbMutex.Lock()
@ -334,7 +337,7 @@ func (s *SqliteBackend) handleJob(ctx context.Context, jobID string) (err error)
return
}
ctx = context.WithValue(ctx, txCtxVarKey, tx)
fmt.Printf("Updated ctx with tx : %+v", ctx)
log.Printf("Updated ctx with tx : %+v", ctx)
err = s.updateJob(ctx, jobErr, internal.JobStatusInProgress)
if err != nil {
if errors.Is(err, context.Canceled) {
@ -419,7 +422,7 @@ func withJobContext(ctx context.Context, j *jobs.Job) context.Context {
}
func (s *SqliteBackend) updateJob(ctx context.Context, jobErr error, status string) (err error) {
fmt.Println("PS::updateJob, job status:", status)
log.Println("PS::updateJob, job status:", status)
errMsg := ""
@ -440,7 +443,7 @@ func (s *SqliteBackend) updateJob(ctx context.Context, jobErr error, status stri
// In progress job
if len(status) != 0 {
fmt.Println("PS::updating job status to in progress")
log.Println("PS::updating job status to in progress")
qstr := "UPDATE neoq_jobs SET status = $1, error = $2 WHERE id = $3"
_, err = tx.ExecContext(ctx, qstr, status, errMsg, job.ID)
return