Files
2026-09-16 06:26:16 +02:00

101 lines
2.8 KiB
Go

package outbox
import (
"context"
"encoding/json"
"errors"
"path/filepath"
"sync"
"testing"
"time"
)
func TestPersistenceDeduplicationAndLeaseRecovery(t *testing.T) {
ctx := context.Background()
path := filepath.Join(t.TempDir(), "outbox.db")
s, err := Open(path)
if err != nil {
t.Fatal(err)
}
jobs := []Job{{MappingID: "one", Target: "webhook", Data: json.RawMessage(`{"payload":"test"}`)}, {MappingID: "two", Target: "smtp", Data: json.RawMessage(`{}`)}}
receipt, err := s.Enqueue(ctx, "source", "key", "digest", jobs)
if err != nil {
t.Fatal(err)
}
duplicate, err := s.Enqueue(ctx, "source", "key", "digest", jobs)
if err != nil || !duplicate.Duplicate || duplicate.ID != receipt.ID {
t.Fatalf("duplicate=%+v err=%v", duplicate, err)
}
if _, err := s.Enqueue(ctx, "source", "key", "different", jobs); !errors.Is(err, ErrConflict) {
t.Fatalf("conflict=%v", err)
}
j, err := s.Claim(ctx, time.Now())
if err != nil || j == nil {
t.Fatalf("claim: %v", err)
}
if err := s.Finish(ctx, *j, "succeeded", 204, "", time.Now()); err != nil {
t.Fatal(err)
}
if err := s.Retry(ctx, j.ID); !errors.Is(err, ErrNotRetryable) {
t.Fatal("successful delivery was retryable")
}
orphan, err := s.Claim(ctx, time.Now())
if err != nil || orphan == nil {
t.Fatal(err)
}
s.Close()
s, err = Open(path)
if err != nil {
t.Fatal(err)
}
defer s.Close()
if j, err := s.Claim(ctx, time.Now()); err != nil || j != nil {
t.Fatalf("active lease should remain protected: %+v %v", j, err)
}
recovered, err := s.Claim(ctx, time.Now().Add(11*time.Minute))
if err != nil || recovered == nil || recovered.ID != orphan.ID || recovered.Attempts != 2 {
t.Fatalf("recovery=%+v %v", recovered, err)
}
if err := s.Finish(ctx, *orphan, "succeeded", 200, "", time.Now()); err == nil {
t.Fatal("stale worker completed a new lease")
}
if err := s.Finish(ctx, *recovered, "dead", 400, "bad request", time.Now()); err != nil {
t.Fatal(err)
}
if err := s.Retry(ctx, recovered.ID); err != nil {
t.Fatal(err)
}
history, err := s.History(ctx, recovered.ID)
if err != nil || len(history) != 2 || history[0].State != "manual_retry" {
t.Fatalf("history=%+v %v", history, err)
}
all, err := s.List(ctx, 50, 0)
if err != nil || len(all) != 2 {
t.Fatalf("jobs=%+v %v", all, err)
}
}
func TestConcurrentAcceptanceIsAtomic(t *testing.T) {
s, err := Open(filepath.Join(t.TempDir(), "outbox.db"))
if err != nil {
t.Fatal(err)
}
defer s.Close()
var wg sync.WaitGroup
for i := 0; i < 12; i++ {
wg.Add(1)
go func() {
defer wg.Done()
_, err := s.Enqueue(context.Background(), "scope", "one", "digest", []Job{{MappingID: "a", Target: "webhook", Data: json.RawMessage(`{}`)}})
if err != nil {
t.Error(err)
}
}()
}
wg.Wait()
rows, err := s.List(context.Background(), 50, 0)
if err != nil || len(rows) != 1 {
t.Fatalf("rows=%+v %v", rows, err)
}
}