Files
yovision/Sense/internal/auditrelay/worker.go
T
QiuSW 4be2386421
Harness governance / validate (push) Has been cancelled
Harness governance / validate (pull_request) Has been cancelled
feat: deliver Sense audits to Bell (T-016)
2026-08-11 00:24:32 +08:00

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:
}
}
}