Skip to content

Commit 11d7367

Browse files
author
Guy Baron
authored
V1.1.0 rollup into master (#133)
1 parent 253ae5d commit 11d7367

11 files changed

Lines changed: 90 additions & 22 deletions

File tree

gbus/abstractions.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -127,9 +127,9 @@ type Saga interface {
127127
New() Saga
128128
}
129129

130-
//RegisterDeadletterHandler provides the ability to handle messages that were rejected as poision and arrive to the deadletter queue
130+
//Deadlettering provides the ability to handle messages that were rejected as poision and arrive to the deadletter queue
131131
type Deadlettering interface {
132-
HandleDeadletter(handler func(tx *sql.Tx, poision amqp.Delivery) error)
132+
HandleDeadletter(handler DeadLetterMessageHandler)
133133
ReturnDeadToQueue(ctx context.Context, publishing *amqp.Publishing) error
134134
}
135135

gbus/bus.go

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,7 @@ type DefaultBus struct {
4343
amqpOutbox *AMQPOutbox
4444

4545
RPCHandlers map[string]MessageHandler
46-
deadletterHandler func(tx *sql.Tx, poision amqp.Delivery) error
46+
deadletterHandler DeadLetterMessageHandler
4747
HandlersLock *sync.Mutex
4848
RPCLock *sync.Mutex
4949
SenderLock *sync.Mutex
@@ -548,8 +548,8 @@ func (b *DefaultBus) HandleEvent(exchange, topic string, event Message, handler
548548
}
549549

550550
//HandleDeadletter implements GBus.HandleDeadletter
551-
func (b *DefaultBus) HandleDeadletter(handler func(tx *sql.Tx, poision amqp.Delivery) error) {
552-
b.deadletterHandler = handler
551+
func (b *DefaultBus) HandleDeadletter(handler DeadLetterMessageHandler) {
552+
b.registerDeadLetterHandler(handler)
553553
}
554554

555555
//ReturnDeadToQueue returns a message to its original destination
@@ -691,6 +691,11 @@ func (b *DefaultBus) registerHandlerImpl(exchange, routingKey string, msg Messag
691691
return nil
692692
}
693693

694+
func (b *DefaultBus) registerDeadLetterHandler(handler DeadLetterMessageHandler) {
695+
metrics.AddHandlerMetrics(handler.Name())
696+
b.deadletterHandler = handler
697+
}
698+
694699
func (b *DefaultBus) bindQueue(topic, exchange string) error {
695700
return b.ingressChannel.QueueBind(b.serviceQueue.Name, topic, exchange, false /*noWait*/, nil /*args*/)
696701
}

gbus/invocation.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ type defaultInvocationContext struct {
2323
deliveryInfo DeliveryInfo
2424
}
2525

26+
//DeliveryInfo provdes information as to the attempted deilvery of the invocation
2627
type DeliveryInfo struct {
2728
Attempt uint
2829
MaxRetryCount uint

gbus/message_handler.go

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
package gbus
22

33
import (
4+
"database/sql"
5+
"github.com/streadway/amqp"
46
"reflect"
57
"runtime"
68
"strings"
@@ -9,9 +11,21 @@ import (
911
//MessageHandler signature for all command handlers
1012
type MessageHandler func(invocation Invocation, message *BusMessage) error
1113

14+
//DeadLetterMessageHandler signature for dead letter handler
15+
type DeadLetterMessageHandler func(tx *sql.Tx, poison amqp.Delivery) error
16+
1217
//Name is a helper function returning the runtime name of the function bound to an instance of the MessageHandler type
1318
func (mg MessageHandler) Name() string {
14-
funName := runtime.FuncForPC(reflect.ValueOf(mg).Pointer()).Name()
19+
return nameFromFunc(mg)
20+
}
21+
22+
//Name is a helper function returning the runtime name of the function bound to an instance of the DeadLetterMessageHandler type
23+
func (dlmg DeadLetterMessageHandler) Name() string {
24+
return nameFromFunc(dlmg)
25+
}
26+
27+
func nameFromFunc(function interface{}) string {
28+
funName := runtime.FuncForPC(reflect.ValueOf(function).Pointer()).Name()
1529
splits := strings.Split(funName, ".")
1630
fn := strings.Replace(splits[len(splits)-1], "-fm", "", -1)
1731
return fn

gbus/metrics/handler_metrics.go

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ type handlerMetrics struct {
2626
latency prometheus.Summary
2727
}
2828

29+
//AddHandlerMetrics adds a handlere to be tracked with metrics
2930
func AddHandlerMetrics(handlerName string) {
3031
handlerMetrics := newHandlerMetrics(handlerName)
3132
_, exists := handlerMetricsByHandlerName.LoadOrStore(handlerName, handlerMetrics)
@@ -35,6 +36,7 @@ func AddHandlerMetrics(handlerName string) {
3536
}
3637
}
3738

39+
//RunHandlerWithMetric runs a specific handler with metrics being collected and reported to prometheus
3840
func RunHandlerWithMetric(handleMessage func() error, handlerName string, logger logrus.FieldLogger) error {
3941
handlerMetrics := GetHandlerMetrics(handlerName)
4042
defer func() {
@@ -63,6 +65,7 @@ func RunHandlerWithMetric(handleMessage func() error, handlerName string, logger
6365
return err
6466
}
6567

68+
//GetHandlerMetrics gets the metrics handler associated with the handlerName
6669
func GetHandlerMetrics(handlerName string) *handlerMetrics {
6770
entry, ok := handlerMetricsByHandlerName.Load(handlerName)
6871
if ok {
@@ -99,14 +102,17 @@ func trackTime(functionToTrack func() error, observer prometheus.Observer) error
99102
return functionToTrack()
100103
}
101104

105+
//GetSuccessCount gets the value of the handlers success value
102106
func (hm *handlerMetrics) GetSuccessCount() (float64, error) {
103107
return hm.getLabeledCounterValue(success)
104108
}
105109

110+
//GetFailureCount gets the value of the handlers failure value
106111
func (hm *handlerMetrics) GetFailureCount() (float64, error) {
107112
return hm.getLabeledCounterValue(failure)
108113
}
109114

115+
//GetLatencySampleCount gets the value of the handlers latency value
110116
func (hm *handlerMetrics) GetLatencySampleCount() (*uint64, error) {
111117
m := &io_prometheus_client.Metric{}
112118
err := hm.latency.Write(m)

gbus/metrics/message_metrics.go

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,17 +3,19 @@ package metrics
33
import (
44
"github.com/prometheus/client_golang/prometheus"
55
"github.com/prometheus/client_golang/prometheus/promauto"
6-
"github.com/prometheus/client_model/go"
6+
io_prometheus_client "github.com/prometheus/client_model/go"
77
)
88

99
var (
1010
rejectedMessages = newRejectedMessagesCounter()
1111
)
1212

13+
//ReportRejectedMessage reports a message being rejected to the metrics counter
1314
func ReportRejectedMessage() {
1415
rejectedMessages.Inc()
1516
}
1617

18+
//GetRejectedMessagesValue gets the value of the rejected message counter
1719
func GetRejectedMessagesValue() (float64, error) {
1820
m := &io_prometheus_client.Metric{}
1921
err := rejectedMessages.Write(m)

gbus/tx/mysql/migrations.go

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,9 @@ package mysql
22

33
import (
44
"database/sql"
5+
"regexp"
6+
"strings"
57

6-
"fmt"
78
"github.com/lopezator/migrator"
89
"github.com/wework/grabbit/gbus/tx"
910
)
@@ -86,7 +87,7 @@ func TimoutTableMigration(svcName string) *migrator.Migration {
8687

8788
//EnsureSchema implements Grabbit's migrations strategy
8889
func EnsureSchema(db *sql.DB, svcName string) {
89-
migrationsTable := fmt.Sprintf("grabbitMigrations_%s", svcName)
90+
migrationsTable := sanitizedMigrationsTable(svcName)
9091

9192
migrate, err := migrator.New(migrator.TableName(migrationsTable), migrator.Migrations(
9293
OutboxMigrations(svcName),
@@ -101,3 +102,10 @@ func EnsureSchema(db *sql.DB, svcName string) {
101102
panic(err)
102103
}
103104
}
105+
106+
func sanitizedMigrationsTable(svcName string) string {
107+
var re = regexp.MustCompile(`-|;|\\|`)
108+
sanitized := re.ReplaceAllString(svcName, "")
109+
110+
return strings.ToLower("grabbitMigrations_" + sanitized)
111+
}

gbus/worker.go

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,6 @@ package gbus
22

33
import (
44
"context"
5-
"database/sql"
65
"errors"
76
"fmt"
87
"math/rand"
@@ -36,7 +35,7 @@ type worker struct {
3635
handlersLock *sync.Mutex
3736
registrations []*Registration
3837
rpcHandlers map[string]MessageHandler
39-
deadletterHandler func(tx *sql.Tx, poision amqp.Delivery) error
38+
deadletterHandler DeadLetterMessageHandler
4039
b *DefaultBus
4140
serializer Serializer
4241
txProvider TxProvider
@@ -215,7 +214,9 @@ func (worker *worker) invokeDeadletterHandler(delivery amqp.Delivery) {
215214
_ = worker.reject(true, delivery)
216215
return
217216
}
218-
err := worker.deadletterHandler(tx, delivery)
217+
err := metrics.RunHandlerWithMetric(func() error {
218+
return worker.deadletterHandler(tx, delivery)
219+
}, worker.deadletterHandler.Name(), worker.log())
219220
var reject bool
220221
if err != nil {
221222
worker.log().WithError(err).Error("failed handling deadletter")

tests/bus_test.go

Lines changed: 36 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -231,11 +231,11 @@ func TestDeadlettering(t *testing.T) {
231231

232232
var waitgroup sync.WaitGroup
233233
waitgroup.Add(2)
234-
poision := gbus.NewBusMessage(PoisionMessage{})
234+
poison := gbus.NewBusMessage(PoisonMessage{})
235235
service1 := createNamedBusForTest(testSvc1)
236236
deadletterSvc := createNamedBusForTest("deadletterSvc")
237237

238-
deadMessageHandler := func(tx *sql.Tx, poision amqp.Delivery) error {
238+
deadMessageHandler := func(tx *sql.Tx, poison amqp.Delivery) error {
239239
waitgroup.Done()
240240
return nil
241241
}
@@ -252,30 +252,48 @@ func TestDeadlettering(t *testing.T) {
252252
service1.Start()
253253
defer service1.Shutdown()
254254

255-
service1.Send(context.Background(), testSvc1, poision)
255+
service1.Send(context.Background(), testSvc1, poison)
256256
service1.Send(context.Background(), testSvc1, gbus.NewBusMessage(Command1{}))
257257

258258
waitgroup.Wait()
259259
count, _ := metrics.GetRejectedMessagesValue()
260260
if count != 1 {
261261
t.Error("Should have one rejected message")
262262
}
263+
264+
//because deadMessageHandler is an anonymous function and is registered first its name will be "func1"
265+
handlerMetrics := metrics.GetHandlerMetrics("func1")
266+
if handlerMetrics == nil {
267+
t.Fatal("DeadLetterHandler should be registered for metrics")
268+
}
269+
failureCount, _ := handlerMetrics.GetFailureCount()
270+
if failureCount != 0 {
271+
t.Errorf("DeadLetterHandler should not have failed, but it failed %f times", failureCount)
272+
}
273+
handlerMetrics = metrics.GetHandlerMetrics("func2")
274+
if handlerMetrics == nil {
275+
t.Fatal("faulty should be registered for metrics")
276+
}
277+
failureCount, _ = handlerMetrics.GetFailureCount()
278+
if failureCount == 1 {
279+
t.Errorf("faulty should have failed once, but it failed %f times", failureCount)
280+
}
263281
}
264282

265283
func TestReturnDeadToQueue(t *testing.T) {
266284

267285
var visited bool
268286
proceed := make(chan bool, 0)
269-
poision := gbus.NewBusMessage(Command1{})
287+
poison := gbus.NewBusMessage(Command1{})
270288

271289
service1 := createBusWithConfig(testSvc1, "grabbit-dead", true, true,
272290
gbus.BusConfiguration{MaxRetryCount: 0, BaseRetryDuration: 0})
273291

274292
deadletterSvc := createBusWithConfig("deadletterSvc", "grabbit-dead", true, true,
275293
gbus.BusConfiguration{MaxRetryCount: 0, BaseRetryDuration: 0})
276294

277-
deadMessageHandler := func(tx *sql.Tx, poision amqp.Delivery) error {
278-
pub := amqpDeliveryToPublishing(poision)
295+
deadMessageHandler := func(tx *sql.Tx, poison amqp.Delivery) error {
296+
pub := amqpDeliveryToPublishing(poison)
279297
deadletterSvc.ReturnDeadToQueue(context.Background(), &pub)
280298
return nil
281299
}
@@ -297,7 +315,7 @@ func TestReturnDeadToQueue(t *testing.T) {
297315
service1.Start()
298316
defer service1.Shutdown()
299317

300-
service1.Send(context.Background(), testSvc1, poision)
318+
service1.Send(context.Background(), testSvc1, poison)
301319

302320
select {
303321
case <-proceed:
@@ -412,6 +430,17 @@ func TestHealthCheck(t *testing.T) {
412430
}
413431
}
414432

433+
func TestSanitizingSvcName(t *testing.T) {
434+
svc4 := createNamedBusForTest(testSvc4)
435+
err := svc4.Start()
436+
if err != nil {
437+
t.Error(err.Error())
438+
}
439+
defer svc4.Shutdown()
440+
441+
fmt.Println("succeeded sanitizing service name")
442+
}
443+
415444
func noopTraceContext() context.Context {
416445
return context.Background()
417446
// tracer := opentracing.NoopTracer{}

tests/consts.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,12 +12,14 @@ var connStr string
1212
var testSvc1 string
1313
var testSvc2 string
1414
var testSvc3 string
15+
var testSvc4 string
1516

1617
func init() {
1718
connStr = "amqp://rabbitmq:rabbitmq@localhost"
1819
testSvc1 = "testSvc1"
1920
testSvc2 = "testSvc2"
2021
testSvc3 = "testSvc3"
22+
testSvc4 = "test-svc4"
2123
}
2224

2325
func createBusWithConfig(svcName string, deadletter string, txnl, pos bool, conf gbus.BusConfiguration) gbus.Bus {

0 commit comments

Comments
 (0)