This repository has been archived by the owner on Dec 9, 2019. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 6
/
invocation.go
73 lines (59 loc) · 1.97 KB
/
invocation.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
package gbus
import (
"context"
"database/sql"
"time"
"github.com/sirupsen/logrus"
)
type defaultInvocationContext struct {
invocingSvc string
bus *DefaultBus
inboundMsg *BusMessage
tx *sql.Tx
ctx context.Context
exchange string
routingKey string
}
func (dfi *defaultInvocationContext) Reply(ctx context.Context, replyMessage *BusMessage) error {
if dfi.inboundMsg != nil {
replyMessage.CorrelationID = dfi.inboundMsg.ID
replyMessage.SagaCorrelationID = dfi.inboundMsg.SagaID
replyMessage.RPCID = dfi.inboundMsg.RPCID
}
var err error
if dfi.tx != nil {
return dfi.bus.sendWithTx(ctx, dfi.tx, dfi.invocingSvc, replyMessage)
}
if err = dfi.bus.Send(ctx, dfi.invocingSvc, replyMessage); err != nil {
//TODO: add logs?
logrus.WithError(err).Error("could not send reply")
}
return err
}
func (dfi *defaultInvocationContext) Send(ctx context.Context, toService string, command *BusMessage, policies ...MessagePolicy) error {
if dfi.tx != nil {
return dfi.bus.sendWithTx(ctx, dfi.tx, toService, command, policies...)
}
return dfi.bus.Send(ctx, toService, command, policies...)
}
func (dfi *defaultInvocationContext) Publish(ctx context.Context, exchange, topic string, event *BusMessage, policies ...MessagePolicy) error {
if dfi.tx != nil {
return dfi.bus.publishWithTx(ctx, dfi.tx, exchange, topic, event, policies...)
}
return dfi.bus.Publish(ctx, exchange, topic, event, policies...)
}
func (dfi *defaultInvocationContext) RPC(ctx context.Context, service string, request, reply *BusMessage, timeout time.Duration) (*BusMessage, error) {
return dfi.bus.RPC(ctx, service, request, reply, timeout)
}
func (dfi *defaultInvocationContext) Bus() Messaging {
return dfi
}
func (dfi *defaultInvocationContext) Tx() *sql.Tx {
return dfi.tx
}
func (dfi *defaultInvocationContext) Ctx() context.Context {
return dfi.ctx
}
func (dfi *defaultInvocationContext) Routing() (exchange, routingKey string) {
return dfi.exchange, dfi.routingKey
}