This commit is contained in:
@@ -204,26 +204,36 @@ func (w *Worker) processJobs(ctx context.Context) {
|
||||
}
|
||||
|
||||
if err := w.setTenantTx(ctx, tx); err != nil {
|
||||
tx.Rollback()
|
||||
if rbErr := tx.Rollback(); rbErr != nil {
|
||||
w.logger.Error("failed to rollback transaction", "error", rbErr)
|
||||
}
|
||||
w.logger.Error("failed to set tenant", "error", err)
|
||||
return
|
||||
}
|
||||
|
||||
jobID, leaseToken, err := w.fetchNextJobTx(ctx, tx)
|
||||
if err != nil {
|
||||
tx.Rollback()
|
||||
if rbErr := tx.Rollback(); rbErr != nil {
|
||||
w.logger.Error("failed to rollback transaction", "error", rbErr)
|
||||
}
|
||||
w.logger.Error("failed to fetch job", "error", err)
|
||||
return
|
||||
}
|
||||
|
||||
if jobID <= 0 {
|
||||
tx.Rollback() // No job found, rollback
|
||||
return // No more jobs
|
||||
// No job found, rollback
|
||||
if rbErr := tx.Rollback(); rbErr != nil {
|
||||
w.logger.Error("failed to rollback transaction", "error", rbErr)
|
||||
}
|
||||
return // No more jobs
|
||||
}
|
||||
|
||||
// Run the job
|
||||
if err := w.runJobTx(ctx, tx, jobID, leaseToken); err != nil {
|
||||
tx.Rollback() // Rollback on genuine infra failure
|
||||
// Rollback on genuine infra failure
|
||||
if rbErr := tx.Rollback(); rbErr != nil {
|
||||
w.logger.Error("failed to rollback transaction", "error", rbErr)
|
||||
}
|
||||
w.logger.Error("failed to run job", "job_id", jobID, "error", err)
|
||||
} else {
|
||||
if err := tx.Commit(); err != nil {
|
||||
@@ -247,13 +257,13 @@ func (w *Worker) setTenantTx(ctx context.Context, tx adapter.DBTransaction) erro
|
||||
|
||||
// fetchNextJobTx fetches the next job from the queue within a transaction,
|
||||
// claiming it with a lease that must be presented back to broker_run.
|
||||
func (w *Worker) fetchNextJobTx(ctx context.Context, tx adapter.DBTransaction) (int64, string, error) {
|
||||
func (w *Worker) fetchNextJobTx(ctx context.Context, tx adapter.DBTransaction) (jobID int64, leaseToken string, err error) {
|
||||
var retval int
|
||||
var errmsg string
|
||||
var nullableJobID sql.NullInt64
|
||||
var nullableLeaseToken sql.NullString
|
||||
|
||||
err := tx.QueryRow(ctx,
|
||||
err = tx.QueryRow(ctx,
|
||||
"SELECT p_retval, p_errmsg, p_job_id, p_lease_token FROM broker.broker_get($1, $2, $3)",
|
||||
w.QueueNumber, w.InstanceID, w.leaseSeconds,
|
||||
).Scan(&retval, &errmsg, &nullableJobID, &nullableLeaseToken)
|
||||
|
||||
Reference in New Issue
Block a user