From 1bd85ad6c50e34655c00c4ed8f4e3c161abbccd0 Mon Sep 17 00:00:00 2001 From: David Thorpe Date: Wed, 17 Jun 2026 17:20:24 +0200 Subject: [PATCH 1/4] Added permanent failure support --- pgqueue/manager/exec.go | 32 +++++++++++++++ pgqueue/manager/exec_test.go | 36 ++++++++++++++++ pgqueue/manager/manager.go | 6 +++ pgqueue/manager/queue_test.go | 77 +++++++++++++++++++++++++++++++++++ pgqueue/schema/objects.sql | 8 +++- 5 files changed, 158 insertions(+), 1 deletion(-) diff --git a/pgqueue/manager/exec.go b/pgqueue/manager/exec.go index 7fb2514..278acc0 100644 --- a/pgqueue/manager/exec.go +++ b/pgqueue/manager/exec.go @@ -189,6 +189,38 @@ func (exec *exec) RunQueueTask(ctx context.Context, task *schema.Task, result ch } func run(ctx context.Context, fn schema.TaskFunc, payload json.RawMessage) (resp *Result) { + return runWithGrace(ctx, fn, payload, time.Minute) +} + +func runWithGrace(ctx context.Context, fn schema.TaskFunc, payload json.RawMessage, grace time.Duration) (resp *Result) { + result := make(chan *Result, 1) + go func() { + result <- runCallback(ctx, fn, payload) + }() + + select { + case resp := <-result: + return resp + case <-ctx.Done(): + // TTL is cooperative via context cancellation. Give callbacks an + // additional grace window to exit before force-failing the task. + if grace <= 0 { + return types.Ptr(Result{Error: context.DeadlineExceeded}) + } + + timer := time.NewTimer(grace) + defer timer.Stop() + + select { + case resp := <-result: + return resp + case <-timer.C: + return types.Ptr(Result{Error: fmt.Errorf("task exceeded TTL and did not stop after %s grace period: %w", grace, context.DeadlineExceeded)}) + } + } +} + +func runCallback(ctx context.Context, fn schema.TaskFunc, payload json.RawMessage) (resp *Result) { defer func() { if recovered := recover(); recovered != nil { var err error diff --git a/pgqueue/manager/exec_test.go b/pgqueue/manager/exec_test.go index 2ab43eb..f934596 100644 --- a/pgqueue/manager/exec_test.go +++ b/pgqueue/manager/exec_test.go @@ -154,3 +154,39 @@ func TestRunQueueTaskRequiresTaskDeadline(t *testing.T) { exec.Close() } + +func TestRunReturnsOnContextDeadlineWhenCallbackBlocks(t *testing.T) { + unblock := make(chan struct{}) + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + + start := time.Now() + result := runWithGrace(ctx, func(context.Context, json.RawMessage) (any, error) { + <-unblock + return map[string]bool{"ok": true}, nil + }, nil, 20*time.Millisecond) + elapsed := time.Since(start) + close(unblock) + + require.NotNil(t, result) + require.Error(t, result.Error) + assert.True(t, errors.Is(result.Error, context.DeadlineExceeded)) + assert.Less(t, elapsed, 500*time.Millisecond) +} + +func TestRunReturnsCallbackErrorIfItExitsDuringGrace(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + + start := time.Now() + result := runWithGrace(ctx, func(ctx context.Context, _ json.RawMessage) (any, error) { + <-ctx.Done() + return nil, ctx.Err() + }, nil, 100*time.Millisecond) + elapsed := time.Since(start) + + require.NotNil(t, result) + require.Error(t, result.Error) + assert.True(t, errors.Is(result.Error, context.DeadlineExceeded)) + assert.Less(t, elapsed, 500*time.Millisecond) +} diff --git a/pgqueue/manager/manager.go b/pgqueue/manager/manager.go index 8923e39..60dfc89 100644 --- a/pgqueue/manager/manager.go +++ b/pgqueue/manager/manager.go @@ -71,6 +71,12 @@ func New(ctx context.Context, pool pg.PoolConn, opts ...Opt) (*Manager, error) { self.PoolConn = pool } + // Ensure at least one task partition exists before any task inserts happen. + if _, err := self.CreateNextPartition(bootstrapCtx); err != nil { + endBootstrapSpan(err) + return nil, err + } + // Register a maintenance ticker if _, err := self.RegisterTicker(bootstrapCtx, schema.DefaultMaintenanceTickerName, schema.TickerMeta{ Interval: types.Ptr(schema.DefaultMaintenancePeriod), diff --git a/pgqueue/manager/queue_test.go b/pgqueue/manager/queue_test.go index 1461d56..eb7dc6a 100644 --- a/pgqueue/manager/queue_test.go +++ b/pgqueue/manager/queue_test.go @@ -5,6 +5,7 @@ import ( "encoding/json" "errors" "sync" + "sync/atomic" "testing" "time" @@ -324,6 +325,82 @@ func TestNextTaskWithZeroConcurrencyKeepsUnlimitedBehavior(t *testing.T) { releaseOnce.Do(func() { close(release) }) } +func TestNextTaskReclaimsExpiredRunningTask(t *testing.T) { + mgr, ctx := test.Begin(t) + defer test.End(t) + + beforeSeq, err := mgr.GetPartitionSeq(ctx) + require.NoError(t, err) + + nextID := beforeSeq + 1 + partitions, err := mgr.ListPartitions(ctx) + require.NoError(t, err) + + createdPartition := "" + if !partitionContains(partitions, nextID) { + meta := schema.PartitionMeta{ + Partition: "task_partition_queue_reclaim_expired", + Start: 1, + End: 1000000, + } + require.NoError(t, mgr.DeletePartition(ctx, meta.Partition)) + require.NoError(t, mgr.CreatePartition(ctx, meta)) + createdPartition = meta.Partition + } + + ttl := 100 * time.Millisecond + concurrency := uint64(1) + release := make(chan struct{}) + var releaseOnce sync.Once + t.Cleanup(func() { + releaseOnce.Do(func() { close(release) }) + }) + started := make(chan struct{}, 2) + var calls atomic.Uint32 + + queue, err := mgr.RegisterQueue(ctx, "queue_reclaim_expired", schema.QueueMeta{ + TTL: &ttl, + Concurrency: &concurrency, + }, func(context.Context, json.RawMessage) (any, error) { + started <- struct{}{} + if calls.Add(1) == 1 { + <-release + } + return nil, nil + }) + require.NoError(t, err) + defer func() { + releaseOnce.Do(func() { close(release) }) + _, _ = mgr.DeleteQueue(ctx, queue.Queue) + if createdPartition != "" { + _ = mgr.DeletePartition(ctx, createdPartition) + } + }() + + _, err = mgr.CreateTask(ctx, queue.Queue, schema.TaskMeta{Payload: json.RawMessage(`{"task":1}`)}) + require.NoError(t, err) + _, err = mgr.CreateTask(ctx, queue.Queue, schema.TaskMeta{Payload: json.RawMessage(`{"task":2}`)}) + require.NoError(t, err) + + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("timed out waiting for first queue task to start") + } + + time.Sleep(ttl + 75*time.Millisecond) + _, err = mgr.CreateTask(ctx, queue.Queue, schema.TaskMeta{Payload: json.RawMessage(`{"task":3}`)}) + require.NoError(t, err) + + // With the fix, the expired running task no longer blocks queue progress, + // so another task attempt can start after TTL expiry. + select { + case <-started: + case <-time.After(2 * time.Second): + t.Fatal("timed out waiting for queue progress after expired running task") + } +} + func TestRunQueueTaskResultIsReleasedDone(t *testing.T) { mgr, ctx := test.Begin(t) defer test.End(t) diff --git a/pgqueue/schema/objects.sql b/pgqueue/schema/objects.sql index de7e5cf..1f8b826 100644 --- a/pgqueue/schema/objects.sql +++ b/pgqueue/schema/objects.sql @@ -170,6 +170,8 @@ WITH selected AS ( "started_at" IS NOT NULL AND "finished_at" IS NULL + AND + ("dies_at" IS NULL OR "dies_at" >= NOW()) GROUP BY "queue" ) retained @@ -178,7 +180,11 @@ WITH selected AS ( WHERE (CARDINALITY(q) = 0 OR t."queue" = ANY(q)) AND - (t."started_at" IS NULL AND t."finished_at" IS NULL) + ( + (t."started_at" IS NULL AND t."finished_at" IS NULL) + OR + (t."started_at" IS NOT NULL AND t."finished_at" IS NULL AND t."dies_at" IS NOT NULL AND t."dies_at" < NOW()) + ) AND (t."delayed_at" IS NULL OR t."delayed_at" <= NOW()) AND From 8bf79f10e6adcf3dca4f6092fd8d235de2bf9b23 Mon Sep 17 00:00:00 2001 From: David Thorpe Date: Wed, 17 Jun 2026 17:37:08 +0200 Subject: [PATCH 2/4] Updated --- pgqueue/manager/run.go | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/pgqueue/manager/run.go b/pgqueue/manager/run.go index 9c7ff8c..8ab5ca1 100644 --- a/pgqueue/manager/run.go +++ b/pgqueue/manager/run.go @@ -3,6 +3,7 @@ package manager import ( "context" "encoding/json" + "errors" "log/slog" "runtime" "strings" @@ -169,7 +170,12 @@ func (manager *Manager) Run(ctx context.Context, log *slog.Logger) error { } status := "" if _, err := manager.ReleaseTask(resultCtx, result.Task.Id, success, releaseResult, &status); err != nil { - log.ErrorContext(resultCtx, "ReleaseTask failed", "queue", result.Queue, "task", result.Task.Id, "error", err.Error()) + if errors.Is(err, pg.ErrNotFound) { + log.WarnContext(resultCtx, "Late queue task result ignored", "queue", result.Queue, "task", result.Task.Id) + queueTimer.Reset(0) + } else { + log.ErrorContext(resultCtx, "ReleaseTask failed", "queue", result.Queue, "task", result.Task.Id, "error", err.Error()) + } continue } if result.Error != nil { From c1d6d0471fe8ea84a122140c7c75fab6742d3c54 Mon Sep 17 00:00:00 2001 From: David Thorpe Date: Wed, 17 Jun 2026 17:39:59 +0200 Subject: [PATCH 3/4] Updates after PR reviews --- pgqueue/manager/exec.go | 9 +++++++-- pgqueue/manager/exec_test.go | 23 +++++++++++++++++++++++ pgqueue/manager/queue_test.go | 16 ++++++++++++++-- 3 files changed, 44 insertions(+), 4 deletions(-) diff --git a/pgqueue/manager/exec.go b/pgqueue/manager/exec.go index 278acc0..8868fbf 100644 --- a/pgqueue/manager/exec.go +++ b/pgqueue/manager/exec.go @@ -202,10 +202,15 @@ func runWithGrace(ctx context.Context, fn schema.TaskFunc, payload json.RawMessa case resp := <-result: return resp case <-ctx.Done(): + cancelErr := ctx.Err() + if cancelErr == nil { + cancelErr = context.Canceled + } + // TTL is cooperative via context cancellation. Give callbacks an // additional grace window to exit before force-failing the task. if grace <= 0 { - return types.Ptr(Result{Error: context.DeadlineExceeded}) + return types.Ptr(Result{Error: cancelErr}) } timer := time.NewTimer(grace) @@ -215,7 +220,7 @@ func runWithGrace(ctx context.Context, fn schema.TaskFunc, payload json.RawMessa case resp := <-result: return resp case <-timer.C: - return types.Ptr(Result{Error: fmt.Errorf("task exceeded TTL and did not stop after %s grace period: %w", grace, context.DeadlineExceeded)}) + return types.Ptr(Result{Error: fmt.Errorf("task did not stop after %s grace period following context cancellation: %w", grace, cancelErr)}) } } } diff --git a/pgqueue/manager/exec_test.go b/pgqueue/manager/exec_test.go index f934596..13ea721 100644 --- a/pgqueue/manager/exec_test.go +++ b/pgqueue/manager/exec_test.go @@ -190,3 +190,26 @@ func TestRunReturnsCallbackErrorIfItExitsDuringGrace(t *testing.T) { assert.True(t, errors.Is(result.Error, context.DeadlineExceeded)) assert.Less(t, elapsed, 500*time.Millisecond) } + +func TestRunWithGracePreservesCanceledCause(t *testing.T) { + unblock := make(chan struct{}) + ctx, cancel := context.WithCancel(context.Background()) + + start := time.Now() + go func() { + time.Sleep(25 * time.Millisecond) + cancel() + }() + + result := runWithGrace(ctx, func(context.Context, json.RawMessage) (any, error) { + <-unblock + return map[string]bool{"ok": true}, nil + }, nil, 20*time.Millisecond) + elapsed := time.Since(start) + close(unblock) + + require.NotNil(t, result) + require.Error(t, result.Error) + assert.True(t, errors.Is(result.Error, context.Canceled)) + assert.Less(t, elapsed, 500*time.Millisecond) +} diff --git a/pgqueue/manager/queue_test.go b/pgqueue/manager/queue_test.go index eb7dc6a..d5fba76 100644 --- a/pgqueue/manager/queue_test.go +++ b/pgqueue/manager/queue_test.go @@ -377,7 +377,7 @@ func TestNextTaskReclaimsExpiredRunningTask(t *testing.T) { } }() - _, err = mgr.CreateTask(ctx, queue.Queue, schema.TaskMeta{Payload: json.RawMessage(`{"task":1}`)}) + firstTask, err := mgr.CreateTask(ctx, queue.Queue, schema.TaskMeta{Payload: json.RawMessage(`{"task":1}`)}) require.NoError(t, err) _, err = mgr.CreateTask(ctx, queue.Queue, schema.TaskMeta{Payload: json.RawMessage(`{"task":2}`)}) require.NoError(t, err) @@ -388,7 +388,19 @@ func TestNextTaskReclaimsExpiredRunningTask(t *testing.T) { t.Fatal("timed out waiting for first queue task to start") } - time.Sleep(ttl + 75*time.Millisecond) + require.Eventually(t, func() bool { + list, err := mgr.ListTasks(ctx, schema.TaskListRequest{Status: "expired"}) + if err != nil { + return false + } + for _, item := range list.Body { + if item.Id == firstTask.Id { + return true + } + } + return false + }, 3*time.Second, 25*time.Millisecond, "timed out waiting for first task to expire") + _, err = mgr.CreateTask(ctx, queue.Queue, schema.TaskMeta{Payload: json.RawMessage(`{"task":3}`)}) require.NoError(t, err) From 1f0fde8637b124e2bdb271c91923a2f916f11629 Mon Sep 17 00:00:00 2001 From: David Thorpe Date: Wed, 17 Jun 2026 17:40:57 +0200 Subject: [PATCH 4/4] Removed logging for null results --- pgqueue/manager/run.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pgqueue/manager/run.go b/pgqueue/manager/run.go index 8ab5ca1..dc27329 100644 --- a/pgqueue/manager/run.go +++ b/pgqueue/manager/run.go @@ -1,6 +1,7 @@ package manager import ( + "bytes" "context" "encoding/json" "errors" @@ -194,7 +195,7 @@ func (manager *Manager) Run(ctx context.Context, log *slog.Logger) error { // Completed ticker task if result.Error != nil { log.ErrorContext(resultCtx, "RunTickerTask result failed", "ticker", result.Ticker, "error", result.Error.Error()) - } else { + } else if len(result.Result) > 0 && !bytes.Equal(result.Result, []byte("null")) { log.InfoContext(resultCtx, "RunTickerTask result", "ticker", result.Ticker, "result", result) } }