package service import ( "context" "encoding/json" "testing" "time" "github.com/DATA-DOG/go-sqlmock" "gorm.io/driver/postgres" "gorm.io/gorm" "silk-server-go/internal/model" ) func newOutboxMock(t *testing.T) (*Outbox, sqlmock.Sqlmock) { t.Helper() sqlDB, mock, err := sqlmock.New() if err != nil { t.Fatalf("create sqlmock: %v", err) } gdb, err := gorm.Open(postgres.New(postgres.Config{Conn: sqlDB}), &gorm.Config{}) if err != nil { t.Fatalf("open gorm: %v", err) } o := NewOutbox(gdb) return o, mock } func TestPublishTxUsesSameEventIDIdempotently(t *testing.T) { o, mock := newOutboxMock(t) event := Event{ ID: "event-1", Type: OutboxEventWechatSubscribe, Payload: json.RawMessage(`{"userId":"u1"}`), } mock.ExpectQuery(`INSERT INTO "outbox_events".*ON CONFLICT \("event_id"\) DO NOTHING`). WithArgs( sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), ). WillReturnRows(sqlmock.NewRows([]string{"id", "payload", "next_attempt_at"}). AddRow("00000000-0000-0000-0000-000000000001", json.RawMessage(`{"userId":"u1"}`), time.Now())) if err := o.PublishTx(o.db, event); err != nil { t.Fatalf("first publish failed: %v", err) } // 第二次相同 eventId 仍走 ON CONFLICT DO NOTHING,不产生重复发送。 mock.ExpectQuery(`INSERT INTO "outbox_events".*ON CONFLICT \("event_id"\) DO NOTHING`). WithArgs( sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), sqlmock.AnyArg(), ). WillReturnRows(sqlmock.NewRows([]string{"id", "payload", "next_attempt_at"})) if err := o.PublishTx(o.db, event); err != nil { t.Fatalf("second publish failed: %v", err) } } func TestProcessPendingWithNoEventsIsNoop(t *testing.T) { o, mock := newOutboxMock(t) o.handler = func(ctx context.Context, event model.OutboxEvent) error { return nil } mock.ExpectQuery(`SELECT .* FROM "outbox_events" WHERE status IN \(.*\)`). WillReturnRows(sqlmock.NewRows([]string{ "id", "event_id", "event_type", "aggregate_type", "aggregate_id", "payload", "status", "attempts", "max_attempts", "next_attempt_at", "last_error", "created_at", "updated_at", })) processed, err := o.ProcessPending(context.Background()) if err != nil { t.Fatalf("ProcessPending failed: %v", err) } if processed != 0 { t.Fatalf("processed = %d, want 0", processed) } } func TestClaimNextReturnsPendingAfterRestart(t *testing.T) { o, mock := newOutboxMock(t) now := time.Now() mock.ExpectQuery(`SELECT .* FROM "outbox_events" WHERE status IN \(.*\)`). WithArgs(OutboxStatusPending, OutboxStatusRetry, sqlmock.AnyArg(), sqlmock.AnyArg()). WillReturnRows(sqlmock.NewRows([]string{ "id", "event_id", "event_type", "aggregate_type", "aggregate_id", "payload", "status", "attempts", "max_attempts", "next_attempt_at", "last_error", "created_at", "updated_at", }).AddRow( "id-1", "event-1", OutboxEventWechatSubscribe, nil, nil, json.RawMessage(`{}`), OutboxStatusPending, 0, 5, now, nil, now, now, )) events, err := o.ClaimNext(context.Background(), 1) if err != nil { t.Fatalf("ClaimNext failed: %v", err) } if len(events) != 1 || events[0].EventID != "event-1" { t.Fatalf("events = %+v", events) } }