forked from jjeffery/stomp
-
Notifications
You must be signed in to change notification settings - Fork 92
/
subscription.go
84 lines (70 loc) · 1.94 KB
/
subscription.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
74
75
76
77
78
79
80
81
82
83
84
package client
import (
"github.com/go-stomp/stomp/v3/frame"
)
type Subscription struct {
conn *Conn
dest string
id string // client's subscription id
ack string // auto, client, client-individual
msgId uint64 // message-id (or ack) for acknowledgement
subList *SubscriptionList // am I in a list
frame *frame.Frame // message allocated to subscription
}
func newSubscription(c *Conn, dest string, id string, ack string) *Subscription {
return &Subscription{
conn: c,
dest: dest,
id: id,
ack: ack,
}
}
func (s *Subscription) Destination() string {
return s.dest
}
func (s *Subscription) Ack() string {
return s.ack
}
func (s *Subscription) Id() string {
return s.id
}
func (s *Subscription) IsAckedBy(msgId uint64) bool {
switch s.ack {
case frame.AckAuto:
return true
case frame.AckClient:
// any later message acknowledges an earlier message
return msgId >= s.msgId
case frame.AckClientIndividual:
return msgId == s.msgId
}
// should not get here
panic("invalid value for subscript.ack")
}
func (s *Subscription) IsNackedBy(msgId uint64) bool {
// TODO: not sure about this, interpreting NACK
// to apply to an individual message
return msgId == s.msgId
}
func (s *Subscription) SendQueueFrame(f *frame.Frame) {
s.setSubscriptionHeader(f)
s.frame = f
// let the connection deal with the subscription
// acknowledgement
s.conn.subChannel <- s
}
// Send a message frame to the client, as part of this
// subscription. Called within the queue when a message
// frame is available.
func (s *Subscription) SendTopicFrame(f *frame.Frame) {
s.setSubscriptionHeader(f)
// topics are handled differently, they just go
// straight to the client without acknowledgement
s.conn.writeChannel <- f
}
func (s *Subscription) setSubscriptionHeader(f *frame.Frame) {
if s.frame != nil {
panic("subscription already has a frame pending")
}
f.Header.Set(frame.Subscription, s.id)
}