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) } }