package auditrelay import ( "context" "errors" "fmt" "regexp" "time" ) var stableErrorCode = regexp.MustCompile(`^[a-z][a-z0-9_]{0,63}$`) type Sender interface { Send(context.Context, []Envelope) ([]Result, error) } type Worker struct { repository Repository sender Sender owner string } func NewWorker(repository Repository, sender Sender, owner string) (*Worker, error) { if repository == nil || sender == nil || owner == "" { return nil, errors.New("audit relay worker dependencies are required") } return &Worker{repository: repository, sender: sender, owner: owner}, nil } func RetryDelay(attempt int) time.Duration { if attempt < 1 { attempt = 1 } if attempt > 9 { return 300 * time.Second } delay := time.Second << (attempt - 1) if delay > 300*time.Second { return 300 * time.Second } return delay } func (w *Worker) RelayOnce(ctx context.Context) error { queued, err := w.repository.ClaimAuditRelayBatch(ctx, w.owner, MaxBatchSize, LeaseDuration) if err != nil || len(queued) == 0 { return err } events := make([]Envelope, len(queued)) for index := range queued { events[index] = queued[index].Envelope } results, sendErr := w.sender.Send(ctx, events) if sendErr != nil { completions := make([]Completion, len(queued)) for index, value := range queued { completions[index] = Completion{EventID: value.Event.EventID, LeaseToken: value.LeaseToken, Disposition: Retry, ErrorCode: "delivery_failed", RetryAfter: RetryDelay(value.AttemptCount)} } if err := w.repository.CompleteAuditRelayBatch(ctx, w.owner, completions); err != nil { return err } return sendErr } if len(results) != len(queued) { return w.retryAll(ctx, queued, "invalid_response") } expected := make(map[string]bool, len(queued)) for _, value := range queued { expected[value.Event.EventID] = true } byID := make(map[string]Result, len(results)) for _, result := range results { validStatus := ((result.Status == "accepted" || result.Status == "duplicate") && result.ErrorCode == nil) || (result.Status == "rejected" && result.ErrorCode != nil && stableErrorCode.MatchString(*result.ErrorCode)) if !expected[result.EventID] || !validStatus { return w.retryAll(ctx, queued, "invalid_response") } if _, duplicate := byID[result.EventID]; duplicate { return w.retryAll(ctx, queued, "invalid_response") } byID[result.EventID] = result } completions := make([]Completion, 0, len(queued)) for _, value := range queued { result, ok := byID[value.Event.EventID] if !ok { return w.retryAll(ctx, queued, "invalid_response") } switch result.Status { case "accepted", "duplicate": completions = append(completions, Completion{EventID: value.Event.EventID, LeaseToken: value.LeaseToken, Disposition: Delivered}) case "rejected": completions = append(completions, Completion{EventID: value.Event.EventID, LeaseToken: value.LeaseToken, Disposition: DeadLetter, ErrorCode: *result.ErrorCode}) } } return w.repository.CompleteAuditRelayBatch(ctx, w.owner, completions) } func (w *Worker) retryAll(ctx context.Context, queued []QueuedEvent, code string) error { values := make([]Completion, len(queued)) for index, value := range queued { values[index] = Completion{EventID: value.Event.EventID, LeaseToken: value.LeaseToken, Disposition: Retry, ErrorCode: code, RetryAfter: RetryDelay(value.AttemptCount)} } if err := w.repository.CompleteAuditRelayBatch(ctx, w.owner, values); err != nil { return err } return fmt.Errorf("audit relay %s", code) } func (w *Worker) Run(ctx context.Context, interval time.Duration, report func(error)) { if interval <= 0 { interval = time.Second } for { if err := w.RelayOnce(ctx); err != nil && ctx.Err() == nil && report != nil { report(err) } timer := time.NewTimer(interval) select { case <-ctx.Done(): timer.Stop() return case <-timer.C: } } }