-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathmsg.go
More file actions
165 lines (133 loc) · 3.23 KB
/
Copy pathmsg.go
File metadata and controls
165 lines (133 loc) · 3.23 KB
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
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
package peanats
import (
"context"
"net/textproto"
"time"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/mikluko/peanats/codec"
)
type Header = textproto.MIMEHeader
type Msg interface {
Subject() string
Header() Header
Data() []byte
}
type Metadatable interface {
Metadata() (*jetstream.MsgMetadata, error)
}
type Respondable interface {
Respond(context.Context, any) error
RespondHeader(context.Context, any, Header) error
RespondMsg(context.Context, Msg) error
}
type Ackable interface {
Ack(context.Context) error
Nak(context.Context) error
NackWithDelay(context.Context, time.Duration) error
Term(context.Context) error
TermWithReason(context.Context, string) error
InProgress(context.Context) error
}
type MsgJetstream interface {
Msg
Metadatable
Ackable
}
type MsgHandler interface {
HandleMsg(context.Context, Msg) error
}
type MsgHandlerFunc func(context.Context, Msg) error
func (f MsgHandlerFunc) HandleMsg(ctx context.Context, m Msg) error {
return f(ctx, m)
}
type MsgMiddleware func(MsgHandler) MsgHandler
func ChainMsgMiddleware(h MsgHandler, mw ...MsgMiddleware) MsgHandler {
for i := range mw {
h = mw[i](h)
}
return h
}
func NewMsg(m *nats.Msg) Msg {
return &msgImpl{m, nil}
}
type msgImpl struct {
*nats.Msg
header Header
}
func (m *msgImpl) Subject() string {
return m.Msg.Subject
}
func (m *msgImpl) Header() Header {
if m.header == nil {
m.header = canonicalizeHeader(m.Msg.Header)
}
return m.header
}
func (m *msgImpl) Data() []byte {
return m.Msg.Data
}
func (m *msgImpl) Respond(ctx context.Context, x any) error {
return m.RespondHeader(ctx, x, nil)
}
func (m *msgImpl) RespondHeader(_ context.Context, x any, header Header) error {
if header == nil {
header = make(Header)
}
data, err := codec.MarshalHeader(x, header)
if err != nil {
return err
}
return m.Msg.RespondMsg(&nats.Msg{
Data: data,
Header: nats.Header(header),
})
}
func (m *msgImpl) RespondMsg(_ context.Context, msg Msg) error {
return m.Msg.RespondMsg(&nats.Msg{
Data: msg.Data(),
Header: nats.Header(msg.Header()),
})
}
func NewJetstream(m jetstream.Msg) MsgJetstream {
return &msgJetstreamImpl{m, nil}
}
type msgJetstreamImpl struct {
jetstream.Msg
header Header
}
func (m *msgJetstreamImpl) Ack(_ context.Context) error {
return m.Msg.Ack()
}
func (m *msgJetstreamImpl) Nak(_ context.Context) error {
return m.Msg.Nak()
}
func (m *msgJetstreamImpl) NackWithDelay(_ context.Context, d time.Duration) error {
return m.Msg.NakWithDelay(d)
}
func (m *msgJetstreamImpl) Term(_ context.Context) error {
return m.Msg.Term()
}
func (m *msgJetstreamImpl) TermWithReason(_ context.Context, s string) error {
return m.Msg.TermWithReason(s)
}
func (m *msgJetstreamImpl) InProgress(_ context.Context) error {
return m.Msg.InProgress()
}
func (m *msgJetstreamImpl) Header() Header {
if m.header == nil {
m.header = canonicalizeHeader(m.Msg.Headers())
}
return m.header
}
// canonicalizeHeader converts a nats.Header to a Header with keys canonicalization.
func canonicalizeHeader(h nats.Header) Header {
if h == nil {
return make(Header)
}
canonical := make(Header, len(h))
for k, v := range h {
canonical[textproto.CanonicalMIMEHeaderKey(k)] = v
}
return canonical
}