mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-13 17:52:03 +00:00
feat: add durable outbox dispatcher
This commit is contained in:
@@ -34,6 +34,57 @@ WHERE event_id = $1 AND published_at IS NULL`
|
||||
|
||||
var ErrOutboxEventNotFound = fmt.Errorf("outbox event not found or already published")
|
||||
|
||||
type OutboxDelivery func(context.Context, OutboxEvent) error
|
||||
|
||||
// OutboxDispatcher is the durable-to-transient bridge. Read and Ack are
|
||||
// injectable so ordering can be tested without a live PostgreSQL instance.
|
||||
// Delivery is at-least-once: a crash after delivery and before acknowledgement
|
||||
// leaves the event replayable, while a delivery failure stops the batch.
|
||||
type OutboxDispatcher struct {
|
||||
Read func(context.Context, int) ([]OutboxEvent, error)
|
||||
Ack func(context.Context, string, time.Time) error
|
||||
Deliver OutboxDelivery
|
||||
}
|
||||
|
||||
func NewOutboxDispatcher(db *sql.DB, deliver OutboxDelivery) *OutboxDispatcher {
|
||||
return &OutboxDispatcher{
|
||||
Read: func(ctx context.Context, limit int) ([]OutboxEvent, error) {
|
||||
return ReadUnpublishedOutbox(ctx, db, limit)
|
||||
},
|
||||
Ack: func(ctx context.Context, eventID string, publishedAt time.Time) error {
|
||||
return MarkOutboxPublished(ctx, db, eventID, publishedAt)
|
||||
},
|
||||
Deliver: deliver,
|
||||
}
|
||||
}
|
||||
|
||||
// Dispatch publishes at most limit events in the store's stable order and
|
||||
// returns the number acknowledged. A successful delivery followed by an ack
|
||||
// error intentionally leaves that event replayable.
|
||||
func (d *OutboxDispatcher) Dispatch(ctx context.Context, limit int, publishedAt time.Time) (int, error) {
|
||||
if d == nil || d.Read == nil || d.Ack == nil || d.Deliver == nil || limit < 1 || limit > 1000 || publishedAt.IsZero() {
|
||||
return 0, fmt.Errorf("invalid outbox dispatcher")
|
||||
}
|
||||
events, err := d.Read(ctx, limit)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
acknowledged := 0
|
||||
for _, event := range events {
|
||||
if event.EventID == "" {
|
||||
return acknowledged, fmt.Errorf("outbox event has no ID")
|
||||
}
|
||||
if err := d.Deliver(ctx, event); err != nil {
|
||||
return acknowledged, err
|
||||
}
|
||||
if err := d.Ack(ctx, event.EventID, publishedAt); err != nil {
|
||||
return acknowledged, err
|
||||
}
|
||||
acknowledged++
|
||||
}
|
||||
return acknowledged, nil
|
||||
}
|
||||
|
||||
// ReadUnpublishedOutbox returns a bounded, stable ordered batch. It does not
|
||||
// mark rows before delivery: a worker crash therefore leaves events replayable.
|
||||
func ReadUnpublishedOutbox(ctx context.Context, db *sql.DB, limit int) ([]OutboxEvent, error) {
|
||||
|
||||
@@ -1,6 +1,9 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"reflect"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
@@ -18,6 +21,38 @@ func TestOutboxSQLPreservesReplayableOrderedReadAndPublishAck(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestOutboxDispatcherAcknowledgesOnlyAfterDelivery(t *testing.T) {
|
||||
events := []OutboxEvent{{EventID: "event-1"}, {EventID: "event-2"}}
|
||||
var delivered, acknowledged []string
|
||||
dispatcher := &OutboxDispatcher{
|
||||
Read: func(context.Context, int) ([]OutboxEvent, error) { return events, nil },
|
||||
Deliver: func(_ context.Context, event OutboxEvent) error {
|
||||
delivered = append(delivered, event.EventID)
|
||||
if event.EventID == "event-2" {
|
||||
return errors.New("transient fan-out failure")
|
||||
}
|
||||
return nil
|
||||
},
|
||||
Ack: func(_ context.Context, eventID string, _ time.Time) error {
|
||||
acknowledged = append(acknowledged, eventID)
|
||||
return nil
|
||||
},
|
||||
}
|
||||
count, err := dispatcher.Dispatch(context.Background(), 10, time.Unix(1000, 0))
|
||||
if err == nil || count != 1 {
|
||||
t.Fatalf("dispatch = (%d, %v), want one acknowledged event and an error", count, err)
|
||||
}
|
||||
if !reflect.DeepEqual(delivered, []string{"event-1", "event-2"}) || !reflect.DeepEqual(acknowledged, []string{"event-1"}) {
|
||||
t.Fatalf("delivery/ack order = %v/%v", delivered, acknowledged)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOutboxDispatcherRejectsInvalidConfiguration(t *testing.T) {
|
||||
if count, err := (*OutboxDispatcher)(nil).Dispatch(context.Background(), 1, time.Unix(1000, 0)); err == nil || count != 0 {
|
||||
t.Fatal("nil dispatcher accepted")
|
||||
}
|
||||
}
|
||||
|
||||
func TestOutboxAdaptersRejectUnsafeArgumentsWithoutDatabase(t *testing.T) {
|
||||
if _, err := ReadUnpublishedOutbox(nil, nil, 1); err == nil {
|
||||
t.Fatal("nil database accepted")
|
||||
|
||||
Reference in New Issue
Block a user