Skip to content

Commit a9b90db

Browse files
ajroetkermattayes
andauthored
Add retries to deregister and panic for mailbox (#168)
* Add retries to deregister and panic for mailbox * Make sure we panic when the actor fails to close its mailbox. Co-authored-by: Matt Braymer-Hayes <matt.hayes91@gmail.com>
1 parent 7909a4c commit a9b90db

2 files changed

Lines changed: 32 additions & 12 deletions

File tree

mailbox.go

Lines changed: 25 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -2,9 +2,19 @@ package grid
22

33
import (
44
"context"
5+
"errors"
56
"fmt"
67
"strings"
78
"sync"
9+
"time"
10+
11+
"github.com/lytics/retry"
12+
)
13+
14+
var (
15+
// errDeregisteredFailed is used internally when we can't deregister a key from etcd.
16+
// It's used by Server.startActorC() to ensure we panic.
17+
errDeregisteredFailed = errors.New("grid: deregistered failed")
818
)
919

1020
// Mailbox for receiving messages.
@@ -15,7 +25,7 @@ type Mailbox struct {
1525
C <-chan Request
1626
c chan Request
1727
closed bool
18-
cleanup func() error
28+
cleanup func()
1929
}
2030

2131
// Close the mailbox.
@@ -27,8 +37,9 @@ func (box *Mailbox) Close() error {
2737
box.closed = true
2838
close(box.c)
2939

40+
box.cleanup()
3041
// Run server provided clean up.
31-
return box.cleanup()
42+
return nil
3243
}
3344

3445
// Name of mailbox, without namespace.
@@ -136,7 +147,7 @@ func newMailbox(s *Server, name, nsName string, size int) (*Mailbox, error) {
136147
}
137148

138149
boxC := make(chan Request, size)
139-
cleanup := func() error {
150+
cleanup := func() {
140151
s.mumb.Lock()
141152
defer s.mumb.Unlock()
142153

@@ -145,12 +156,17 @@ func newMailbox(s *Server, name, nsName string, size int) (*Mailbox, error) {
145156
delete(s.mailboxes, nsName)
146157

147158
// Deregister the name.
148-
timeout, cancel := context.WithTimeout(context.Background(), s.cfg.Timeout)
149-
defer cancel()
150-
err := s.registry.Deregister(timeout, nsName)
151-
152-
// Return any error from the deregister call.
153-
return err
159+
var err error
160+
retry.X(3, 3*time.Second,
161+
func() bool {
162+
timeout, cancel := context.WithTimeout(context.Background(), s.cfg.Timeout)
163+
err = s.registry.Deregister(timeout, nsName)
164+
cancel()
165+
return err != nil
166+
})
167+
if err != nil {
168+
panic(fmt.Errorf("%w: unable to deregister mailbox: %v, error: %v", errDeregisteredFailed, nsName, err))
169+
}
154170
}
155171
box := &Mailbox{
156172
name: name,

server.go

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -444,10 +444,15 @@ func (s *Server) startActorC(c context.Context, start *ActorStart) error {
444444
go func() {
445445
defer s.deregisterActor(nsName)
446446
defer func() {
447-
if err := recover(); err != nil {
447+
if r := recover(); r != nil {
448+
if err, ok := r.(error); ok && errors.Is(err, errDeregisteredFailed) {
449+
// NOTE (2021-06) (mh): We need to panic here
450+
panic(err)
451+
}
452+
448453
stack := niceStack(debug.Stack())
449454
s.logf("panic in namespace: %v, actor: %v, recovered from: %v, stack trace: %v",
450-
s.cfg.Namespace, start.Name, err, stack)
455+
s.cfg.Namespace, start.Name, r, stack)
451456
}
452457
}()
453458
actor.Act(actorCtx)
@@ -466,7 +471,6 @@ func (s *Server) deregisterActor(nsName string) {
466471
return err != nil
467472
})
468473
if err != nil {
469-
s.logf("failed to deregister actor: %v, error: %v", nsName, err)
470474
panic(fmt.Sprintf("unable to deregister actor: %v, error: %v", nsName, err))
471475
}
472476
}

0 commit comments

Comments
 (0)