101 lines
2.8 KiB
Go
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)
|
|
}
|
|
}
|