收窄 assignment 队列职责
This commit is contained in:
@@ -48,8 +48,8 @@ func (p Publisher) Dispatch(ctx context.Context, assignment taskassignment.Assig
|
||||
return nil
|
||||
}
|
||||
|
||||
type Handler interface {
|
||||
Handle(context.Context, taskassignment.Assignment) (bool, error)
|
||||
type Accepter interface {
|
||||
Accept(context.Context, taskassignment.Assignment) (bool, error)
|
||||
}
|
||||
|
||||
// Message is the subset of jetstream.Msg needed by one reconciliation.
|
||||
@@ -57,27 +57,25 @@ type Message interface {
|
||||
Data() []byte
|
||||
DoubleAck(context.Context) error
|
||||
NakWithDelay(time.Duration) error
|
||||
InProgress() error
|
||||
TermWithReason(string) error
|
||||
}
|
||||
|
||||
// Processor maps one delivery to one idempotent worker reconciliation.
|
||||
type Processor struct {
|
||||
TrustDomain string
|
||||
Handler Handler
|
||||
RetryDelay time.Duration
|
||||
PollInterval time.Duration
|
||||
TrustDomain string
|
||||
Accepter Accepter
|
||||
RetryDelay time.Duration
|
||||
}
|
||||
|
||||
func (p Processor) Process(ctx context.Context, message Message) error {
|
||||
if p.Handler == nil {
|
||||
return errors.New("assignment handler is required")
|
||||
if p.Accepter == nil {
|
||||
return errors.New("assignment accepter is required")
|
||||
}
|
||||
assignment, err := taskassignment.Unmarshal(message.Data(), p.TrustDomain)
|
||||
if err != nil {
|
||||
return errors.Join(err, message.TermWithReason("invalid assignment"))
|
||||
}
|
||||
done, err := p.Handler.Handle(ctx, assignment)
|
||||
accepted, err := p.Accepter.Accept(ctx, assignment)
|
||||
if err != nil {
|
||||
delay := p.RetryDelay
|
||||
if delay <= 0 {
|
||||
@@ -85,60 +83,17 @@ func (p Processor) Process(ctx context.Context, message Message) error {
|
||||
}
|
||||
return errors.Join(err, message.NakWithDelay(delay))
|
||||
}
|
||||
if done {
|
||||
if accepted {
|
||||
if err := message.DoubleAck(ctx); err != nil {
|
||||
return fmt.Errorf("ack assignment %s: %w", assignment.ID, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if err := message.InProgress(); err != nil {
|
||||
return fmt.Errorf("extend assignment %s acknowledgement: %w", assignment.ID, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ProcessUntilDone holds one durable delivery while repeatedly reconciling
|
||||
// backend state. Cancellation leaves it unacknowledged for another process.
|
||||
func (p Processor) ProcessUntilDone(ctx context.Context, message Message) error {
|
||||
if p.Handler == nil {
|
||||
return errors.New("assignment handler is required")
|
||||
}
|
||||
assignment, err := taskassignment.Unmarshal(message.Data(), p.TrustDomain)
|
||||
if err != nil {
|
||||
return errors.Join(err, message.TermWithReason("invalid assignment"))
|
||||
}
|
||||
interval := p.PollInterval
|
||||
if interval <= 0 {
|
||||
interval = 2 * time.Second
|
||||
}
|
||||
for {
|
||||
done, handleErr := p.Handler.Handle(ctx, assignment)
|
||||
if handleErr != nil {
|
||||
delay := p.RetryDelay
|
||||
if delay <= 0 {
|
||||
delay = 15 * time.Second
|
||||
}
|
||||
return errors.Join(handleErr, message.NakWithDelay(delay))
|
||||
}
|
||||
if done {
|
||||
if err := message.DoubleAck(ctx); err != nil {
|
||||
return fmt.Errorf("ack assignment %s: %w", assignment.ID, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if err := message.InProgress(); err != nil {
|
||||
return fmt.Errorf("extend assignment %s acknowledgement: %w", assignment.ID, err)
|
||||
}
|
||||
timer := time.NewTimer(interval)
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
if !timer.Stop() {
|
||||
<-timer.C
|
||||
}
|
||||
return ctx.Err()
|
||||
case <-timer.C:
|
||||
}
|
||||
delay := p.RetryDelay
|
||||
if delay <= 0 {
|
||||
delay = 2 * time.Second
|
||||
}
|
||||
return message.NakWithDelay(delay)
|
||||
}
|
||||
|
||||
type consumeAPI interface {
|
||||
@@ -198,7 +153,7 @@ func (c ConsumerComponent) Run(ctx context.Context) error {
|
||||
go func() {
|
||||
defer workers.Done()
|
||||
defer func() { <-semaphore }()
|
||||
if err := c.Processor.ProcessUntilDone(ctx, message); err != nil && !errors.Is(err, context.Canceled) && c.OnError != nil {
|
||||
if err := c.Processor.Process(ctx, message); err != nil && !errors.Is(err, context.Canceled) && c.OnError != nil {
|
||||
c.OnError(err)
|
||||
}
|
||||
}()
|
||||
|
||||
@@ -52,32 +52,25 @@ func TestPublisherUsesBackendSubjectAndAssignmentDeduplication(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
type fakeHandler struct {
|
||||
done bool
|
||||
err error
|
||||
remaining int
|
||||
type fakeAccepter struct {
|
||||
accepted bool
|
||||
err error
|
||||
}
|
||||
|
||||
func (h *fakeHandler) Handle(context.Context, taskassignment.Assignment) (bool, error) {
|
||||
if h.remaining > 0 {
|
||||
h.remaining--
|
||||
return false, nil
|
||||
}
|
||||
return h.done, h.err
|
||||
func (a *fakeAccepter) Accept(context.Context, taskassignment.Assignment) (bool, error) {
|
||||
return a.accepted, a.err
|
||||
}
|
||||
|
||||
type fakeMessage struct {
|
||||
data []byte
|
||||
acked int
|
||||
nacked time.Duration
|
||||
inProgress int
|
||||
terminated int
|
||||
}
|
||||
|
||||
func (m *fakeMessage) Data() []byte { return m.data }
|
||||
func (m *fakeMessage) DoubleAck(context.Context) error { m.acked++; return nil }
|
||||
func (m *fakeMessage) NakWithDelay(delay time.Duration) error { m.nacked = delay; return nil }
|
||||
func (m *fakeMessage) InProgress() error { m.inProgress++; return nil }
|
||||
func (m *fakeMessage) TermWithReason(string) error { m.terminated++; return nil }
|
||||
|
||||
func encodedAssignment(t *testing.T) []byte {
|
||||
@@ -89,24 +82,24 @@ func encodedAssignment(t *testing.T) []byte {
|
||||
return data
|
||||
}
|
||||
|
||||
func TestProcessorAcknowledgesOnlyCompletedAssignment(t *testing.T) {
|
||||
func TestProcessorAcknowledgesPersistedHandoff(t *testing.T) {
|
||||
message := &fakeMessage{data: encodedAssignment(t)}
|
||||
processor := Processor{TrustDomain: "ddupan.top", Handler: &fakeHandler{done: true}}
|
||||
processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{accepted: true}}
|
||||
if err := processor.Process(context.Background(), message); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if message.acked != 1 || message.inProgress != 0 || message.nacked != 0 {
|
||||
if message.acked != 1 || message.nacked != 0 {
|
||||
t.Fatalf("message = %#v", message)
|
||||
}
|
||||
}
|
||||
|
||||
func TestProcessorKeepsRunningAssignmentPending(t *testing.T) {
|
||||
func TestProcessorRetriesUntilBackendHandoffIsDurable(t *testing.T) {
|
||||
message := &fakeMessage{data: encodedAssignment(t)}
|
||||
processor := Processor{TrustDomain: "ddupan.top", Handler: &fakeHandler{}}
|
||||
processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{}, RetryDelay: 2 * time.Second}
|
||||
if err := processor.Process(context.Background(), message); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if message.acked != 0 || message.inProgress != 1 {
|
||||
if message.acked != 0 || message.nacked != 2*time.Second {
|
||||
t.Fatalf("message = %#v", message)
|
||||
}
|
||||
}
|
||||
@@ -115,7 +108,7 @@ func TestProcessorRetriesBackendFailureAndTerminatesPoisonMessage(t *testing.T)
|
||||
retry := &fakeMessage{data: encodedAssignment(t)}
|
||||
processor := Processor{
|
||||
TrustDomain: "ddupan.top",
|
||||
Handler: &fakeHandler{err: errors.New("backend unavailable")},
|
||||
Accepter: &fakeAccepter{err: errors.New("backend unavailable")},
|
||||
RetryDelay: time.Minute,
|
||||
}
|
||||
if err := processor.Process(context.Background(), retry); err == nil {
|
||||
@@ -134,38 +127,6 @@ func TestProcessorRetriesBackendFailureAndTerminatesPoisonMessage(t *testing.T)
|
||||
}
|
||||
}
|
||||
|
||||
func TestProcessUntilDoneReconcilesWithoutRedelivery(t *testing.T) {
|
||||
message := &fakeMessage{data: encodedAssignment(t)}
|
||||
processor := Processor{
|
||||
TrustDomain: "ddupan.top",
|
||||
Handler: &fakeHandler{done: true, remaining: 2},
|
||||
PollInterval: time.Millisecond,
|
||||
}
|
||||
if err := processor.ProcessUntilDone(context.Background(), message); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if message.inProgress != 2 || message.acked != 1 || message.nacked != 0 {
|
||||
t.Fatalf("message = %#v", message)
|
||||
}
|
||||
}
|
||||
|
||||
func TestProcessUntilDoneLeavesMessagePendingOnShutdown(t *testing.T) {
|
||||
message := &fakeMessage{data: encodedAssignment(t)}
|
||||
processor := Processor{
|
||||
TrustDomain: "ddupan.top",
|
||||
Handler: &fakeHandler{},
|
||||
PollInterval: time.Hour,
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
if err := processor.ProcessUntilDone(ctx, message); !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("error = %v", err)
|
||||
}
|
||||
if message.acked != 0 || message.nacked != 0 || message.terminated != 0 {
|
||||
t.Fatalf("message = %#v", message)
|
||||
}
|
||||
}
|
||||
|
||||
type fakeConsumerManager struct{ config jetstream.ConsumerConfig }
|
||||
|
||||
func (m *fakeConsumerManager) CreateOrUpdateConsumer(_ context.Context, _ string, config jetstream.ConsumerConfig) (jetstream.Consumer, error) {
|
||||
|
||||
Reference in New Issue
Block a user