127 lines
3.8 KiB
Go
127 lines
3.8 KiB
Go
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:
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|