feat(channels): add Shangwutong connector
This commit is contained in:
@@ -20,6 +20,18 @@ import (
|
||||
|
||||
var ErrWorkerDatabaseRequired = errors.New("worker database is required")
|
||||
|
||||
type permanentError struct{ err error }
|
||||
|
||||
func (e *permanentError) Error() string { return e.err.Error() }
|
||||
func (e *permanentError) Unwrap() error { return e.err }
|
||||
|
||||
func Permanent(err error) error {
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
return &permanentError{err: err}
|
||||
}
|
||||
|
||||
// JobHandler performs one durable background job.
|
||||
type JobHandler func(context.Context, *model.BackgroundJob) error
|
||||
|
||||
@@ -233,13 +245,33 @@ func (wp *WorkerPool) Enqueue(ctx context.Context, jobType string, payload any,
|
||||
if wp.db == nil {
|
||||
return nil, ErrWorkerDatabaseRequired
|
||||
}
|
||||
job, created, err := wp.persistJob(ctx, wp.db, jobType, payload, opts...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if created {
|
||||
wp.Publish(ctx, job)
|
||||
}
|
||||
return job, nil
|
||||
}
|
||||
|
||||
// EnqueueInTransaction persists a job on the caller's transaction. Call
|
||||
// Publish only after the transaction commits successfully.
|
||||
func (wp *WorkerPool) EnqueueInTransaction(ctx context.Context, tx *gorm.DB, jobType string, payload any, opts ...EnqueueOption) (*model.BackgroundJob, bool, error) {
|
||||
if tx == nil {
|
||||
return nil, false, ErrWorkerDatabaseRequired
|
||||
}
|
||||
return wp.persistJob(ctx, tx, jobType, payload, opts...)
|
||||
}
|
||||
|
||||
func (wp *WorkerPool) persistJob(ctx context.Context, db *gorm.DB, jobType string, payload any, opts ...EnqueueOption) (*model.BackgroundJob, bool, error) {
|
||||
if jobType == "" {
|
||||
return nil, errors.New("job type is required")
|
||||
return nil, false, errors.New("job type is required")
|
||||
}
|
||||
|
||||
payloadBytes, err := marshalPayload(payload)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return nil, false, err
|
||||
}
|
||||
|
||||
job := &model.BackgroundJob{
|
||||
@@ -262,25 +294,31 @@ func (wp *WorkerPool) Enqueue(ctx context.Context, jobType string, payload any,
|
||||
|
||||
if job.IdempotencyKey != "" {
|
||||
var existing model.BackgroundJob
|
||||
err := wp.db.WithContext(ctx).Where("idempotency_key = ?", job.IdempotencyKey).First(&existing).Error
|
||||
err := db.WithContext(ctx).Where("idempotency_key = ?", job.IdempotencyKey).First(&existing).Error
|
||||
if err == nil {
|
||||
return &existing, nil
|
||||
return &existing, false, nil
|
||||
}
|
||||
if !errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return nil, err
|
||||
return nil, false, err
|
||||
}
|
||||
}
|
||||
|
||||
if err := wp.db.WithContext(ctx).Create(job).Error; err != nil {
|
||||
if err := db.WithContext(ctx).Create(job).Error; err != nil {
|
||||
if job.IdempotencyKey != "" {
|
||||
var existing model.BackgroundJob
|
||||
if findErr := wp.db.WithContext(ctx).Where("idempotency_key = ?", job.IdempotencyKey).First(&existing).Error; findErr == nil {
|
||||
return &existing, nil
|
||||
if findErr := db.WithContext(ctx).Where("idempotency_key = ?", job.IdempotencyKey).First(&existing).Error; findErr == nil {
|
||||
return &existing, false, nil
|
||||
}
|
||||
}
|
||||
return nil, err
|
||||
return nil, false, err
|
||||
}
|
||||
return job, true, nil
|
||||
}
|
||||
|
||||
func (wp *WorkerPool) Publish(ctx context.Context, job *model.BackgroundJob) {
|
||||
if job == nil {
|
||||
return
|
||||
}
|
||||
// Push to Redis Stream for immediate dispatch if the job is due.
|
||||
// Scheduled (future) jobs are picked up by the sweep goroutine when they mature.
|
||||
if wp.rdb != nil && !job.ScheduledAt.After(wp.now()) {
|
||||
@@ -293,7 +331,6 @@ func (wp *WorkerPool) Enqueue(ctx context.Context, jobType string, payload any,
|
||||
)
|
||||
}
|
||||
}
|
||||
return job, nil
|
||||
}
|
||||
|
||||
func (wp *WorkerPool) Start() error {
|
||||
@@ -676,7 +713,8 @@ func (wp *WorkerPool) fail(ctx context.Context, job *model.BackgroundJob, err er
|
||||
"locked_by": "",
|
||||
"last_error": err.Error(),
|
||||
}
|
||||
if job.Attempts >= job.MaxAttempts {
|
||||
var permanent *permanentError
|
||||
if job.Attempts >= job.MaxAttempts || errors.As(err, &permanent) {
|
||||
updates["status"] = model.BackgroundJobStatusDead
|
||||
updates["failed_at"] = &now
|
||||
} else {
|
||||
|
||||
@@ -136,6 +136,26 @@ func TestWorkerPoolRetriesThenDeadLettersFailures(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestWorkerPoolPermanentErrorDeadLettersImmediately(t *testing.T) {
|
||||
db := newWorkerTestDB(t)
|
||||
wp := NewWorkerPoolWithOptions(db, WithBackoff(func(attempt int) time.Duration { return 0 }))
|
||||
wp.Register("permanent", func(context.Context, *model.BackgroundJob) error {
|
||||
return Permanent(errors.New("invalid webhook payload"))
|
||||
})
|
||||
job, err := wp.Enqueue(context.Background(), "permanent", nil, WithMaxAttempts(10))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
processed, err := wp.ProcessOne(context.Background())
|
||||
if !processed || err == nil {
|
||||
t.Fatalf("processed=%v err=%v", processed, err)
|
||||
}
|
||||
reloaded := loadJob(t, db, job.ID)
|
||||
if reloaded.Status != model.BackgroundJobStatusDead || reloaded.Attempts != 1 {
|
||||
t.Fatalf("expected immediate dead letter, got %+v", reloaded)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWorkerPoolRespectsScheduleAndQueues(t *testing.T) {
|
||||
db := newWorkerTestDB(t)
|
||||
now := time.Date(2026, 6, 5, 10, 0, 0, 0, time.UTC)
|
||||
|
||||
Reference in New Issue
Block a user