Skip to content

Commit 413bb22

Browse files
Fix concurrency security issues in SubscriptionManager and ListenerManager (#186)
fix: Make `SubscriptionManager.Destroy` safe to call concurrently and more than once by guarding teardown with `sync.Once`, flipping `channelsOpen` under the write lock, and closing each exit channel exactly once, eliminating the data race and `close of closed channel` panic. fix: Bound each `ListenerManager.announce*` send with a 60s per-listener deadline so a stalled or absent consumer no longer parks the announce goroutine (and the event it holds) indefinitely, fixing an unbounded goroutine leak; `copyListeners` now takes an `RLock` since it only reads.
1 parent e73b235 commit 413bb22

8 files changed

Lines changed: 341 additions & 91 deletions

.pubnub.yml

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,13 @@
11
---
2-
version: v9.0.3
2+
version: v9.0.4
33
changelog:
4+
- date: 2026-07-13
5+
version: v9.0.4
6+
changes:
7+
- type: bug
8+
text: "Make `SubscriptionManager.Destroy` safe to call concurrently and more than once by guarding teardown with `sync.Once`, flipping `channelsOpen` under the write lock, and closing each exit channel exactly once, eliminating the data race and `close of closed channel` panic."
9+
- type: bug
10+
text: "Bound each `ListenerManager.announce*` send with a 60s per-listener deadline so a stalled or absent consumer no longer parks the announce goroutine (and the event it holds) indefinitely, fixing an unbounded goroutine leak; `copyListeners` now takes an `RLock` since it only reads."
411
- date: 2026-07-02
512
version: v9.0.3
613
changes:
@@ -864,7 +871,7 @@ sdks:
864871
distribution-type: package
865872
distribution-repository: GitHub
866873
package-name: Go
867-
location: https://github.com/pubnub/go/releases/tag/v9.0.3
874+
location: https://github.com/pubnub/go/releases/tag/v9.0.4
868875
requires:
869876
-
870877
name: "Go"

CHANGELOG.md

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,10 @@
1+
## v9.0.4
2+
July 13 2026
3+
4+
#### Fixed
5+
- Make `SubscriptionManager.Destroy` safe to call concurrently and more than once by guarding teardown with `sync.Once`, flipping `channelsOpen` under the write lock, and closing each exit channel exactly once, eliminating the data race and `close of closed channel` panic.
6+
- Bound each `ListenerManager.announce*` send with a 60s per-listener deadline so a stalled or absent consumer no longer parks the announce goroutine (and the event it holds) indefinitely, fixing an unbounded goroutine leak; `copyListeners` now takes an `RLock` since it only reads.
7+
18
## v9.0.3
29
July 02 2026
310

VERSION

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1 +1 @@
1-
9.0.3
1+
9.0.4

listener_manager.go

Lines changed: 58 additions & 74 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,14 @@ package pubnub
22

33
import (
44
"sync"
5+
"time"
56
)
67

8+
// listenerAnnounceTimeout bounds how long an announce goroutine waits to hand an
9+
// event to a listener channel before dropping it, so a stalled consumer cannot
10+
// leak a goroutine (and the event it holds) forever.
11+
const listenerAnnounceTimeout = 60 * time.Second
12+
713
// Listener type has all the `types` of response events
814
type Listener struct {
915
Status chan *PNStatus
@@ -40,6 +46,8 @@ type ListenerManager struct {
4046
exitListener chan bool
4147
exitListenerAnnounce chan bool
4248
pubnub *PubNub
49+
// announceTimeout is the per-listener send deadline for the announce* fan-out.
50+
announceTimeout time.Duration
4351
}
4452

4553
func newListenerManager(ctx Context, pn *PubNub) *ListenerManager {
@@ -49,6 +57,7 @@ func newListenerManager(ctx Context, pn *PubNub) *ListenerManager {
4957
exitListener: make(chan bool),
5058
exitListenerAnnounce: make(chan bool),
5159
pubnub: pn,
60+
announceTimeout: listenerAnnounceTimeout,
5261
}
5362
}
5463

@@ -77,24 +86,50 @@ func (m *ListenerManager) removeAllListeners() {
7786
}
7887

7988
func (m *ListenerManager) copyListeners() map[*Listener]bool {
80-
m.Lock()
81-
lis := make(map[*Listener]bool)
89+
m.RLock()
90+
lis := make(map[*Listener]bool, len(m.listeners))
8291
for k, v := range m.listeners {
8392
lis[k] = v
8493
}
85-
m.Unlock()
94+
m.RUnlock()
8695
return lis
8796
}
8897

98+
// announceEvent delivers one event to a listener channel, returning false when
99+
// the manager is shutting down (exit fired) so the caller stops the fan-out.
100+
// A positive announceTimeout drops the event on a stalled consumer; <= 0 blocks.
101+
func announceEvent[T any](m *ListenerManager, exit <-chan bool, ch chan<- T, event T, source string) bool {
102+
if m.announceTimeout <= 0 {
103+
select {
104+
case <-exit:
105+
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, source+": exit listener", false)
106+
return false
107+
case ch <- event:
108+
return true
109+
}
110+
}
111+
112+
timer := time.NewTimer(m.announceTimeout)
113+
defer timer.Stop()
114+
115+
select {
116+
case <-exit:
117+
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, source+": exit listener", false)
118+
return false
119+
case ch <- event:
120+
return true
121+
case <-timer.C:
122+
m.pubnub.loggerManager.LogSimple(PNLogLevelWarn, source+": slow consumer, dropping event", false)
123+
return true
124+
}
125+
}
126+
89127
func (m *ListenerManager) announceStatus(status *PNStatus) {
90128
go func() {
91129
lis := m.copyListeners()
92130
for l := range lis {
93-
select {
94-
case <-m.exitListener:
95-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announceStatus: exit listener", false)
131+
if !announceEvent(m, m.exitListener, l.Status, status, "announceStatus") {
96132
return
97-
case l.Status <- status:
98133
}
99134
}
100135
}()
@@ -103,31 +138,20 @@ func (m *ListenerManager) announceStatus(status *PNStatus) {
103138
func (m *ListenerManager) announceMessage(message *PNMessage) {
104139
go func() {
105140
lis := m.copyListeners()
106-
AnnounceMessageLabel:
107141
for l := range lis {
108-
select {
109-
case <-m.exitListenerAnnounce:
110-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announceMessage: exit listener", false)
111-
break AnnounceMessageLabel
112-
case l.Message <- message:
142+
if !announceEvent(m, m.exitListenerAnnounce, l.Message, message, "announceMessage") {
143+
return
113144
}
114145
}
115-
116146
}()
117147
}
118148

119149
func (m *ListenerManager) announceSignal(message *PNMessage) {
120150
go func() {
121151
lis := m.copyListeners()
122-
123-
AnnounceSignalLabel:
124152
for l := range lis {
125-
select {
126-
case <-m.exitListener:
127-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announceSignal: exit listener", false)
128-
break AnnounceSignalLabel
129-
130-
case l.Signal <- message:
153+
if !announceEvent(m, m.exitListener, l.Signal, message, "announceSignal") {
154+
return
131155
}
132156
}
133157
}()
@@ -136,16 +160,9 @@ func (m *ListenerManager) announceSignal(message *PNMessage) {
136160
func (m *ListenerManager) announceUUIDEvent(message *PNUUIDEvent) {
137161
go func() {
138162
lis := m.copyListeners()
139-
140-
AnnounceUUIDEventLabel:
141163
for l := range lis {
142-
select {
143-
case <-m.exitListener:
144-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announceUUIDEvent: exit listener", false)
145-
break AnnounceUUIDEventLabel
146-
147-
case l.UUIDEvent <- message:
148-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announceUUIDEvent: message sent", false)
164+
if !announceEvent(m, m.exitListener, l.UUIDEvent, message, "announceUUIDEvent") {
165+
return
149166
}
150167
}
151168
}()
@@ -154,16 +171,9 @@ func (m *ListenerManager) announceUUIDEvent(message *PNUUIDEvent) {
154171
func (m *ListenerManager) announceChannelEvent(message *PNChannelEvent) {
155172
go func() {
156173
lis := m.copyListeners()
157-
158-
AnnounceChannelEventLabel:
159174
for l := range lis {
160-
select {
161-
case <-m.exitListener:
162-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announceChannelEvent: exit listener", false)
163-
break AnnounceChannelEventLabel
164-
165-
case l.ChannelEvent <- message:
166-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announceChannelEvent: message sent", false)
175+
if !announceEvent(m, m.exitListener, l.ChannelEvent, message, "announceChannelEvent") {
176+
return
167177
}
168178
}
169179
}()
@@ -172,16 +182,9 @@ func (m *ListenerManager) announceChannelEvent(message *PNChannelEvent) {
172182
func (m *ListenerManager) announceMembershipEvent(message *PNMembershipEvent) {
173183
go func() {
174184
lis := m.copyListeners()
175-
176-
AnnounceMembershipEvent:
177185
for l := range lis {
178-
select {
179-
case <-m.exitListener:
180-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announceMembershipEvent: exit listener", false)
181-
break AnnounceMembershipEvent
182-
183-
case l.MembershipEvent <- message:
184-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announceMembershipEvent: message sent", false)
186+
if !announceEvent(m, m.exitListener, l.MembershipEvent, message, "announceMembershipEvent") {
187+
return
185188
}
186189
}
187190
}()
@@ -190,16 +193,9 @@ func (m *ListenerManager) announceMembershipEvent(message *PNMembershipEvent) {
190193
func (m *ListenerManager) announceMessageActionsEvent(message *PNMessageActionsEvent) {
191194
go func() {
192195
lis := m.copyListeners()
193-
194-
AnnounceMessageActionsEvent:
195196
for l := range lis {
196-
select {
197-
case <-m.exitListener:
198-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announceMessageActionsEvent: exit listener", false)
199-
break AnnounceMessageActionsEvent
200-
201-
case l.MessageActionsEvent <- message:
202-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announceMessageActionsEvent: message sent", false)
197+
if !announceEvent(m, m.exitListener, l.MessageActionsEvent, message, "announceMessageActionsEvent") {
198+
return
203199
}
204200
}
205201
}()
@@ -208,15 +204,9 @@ func (m *ListenerManager) announceMessageActionsEvent(message *PNMessageActionsE
208204
func (m *ListenerManager) announcePresence(presence *PNPresence) {
209205
go func() {
210206
lis := m.copyListeners()
211-
212-
AnnouncePresenceLabel:
213207
for l := range lis {
214-
select {
215-
case <-m.exitListener:
216-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announcePresence: exit listener", false)
217-
break AnnouncePresenceLabel
218-
219-
case l.Presence <- presence:
208+
if !announceEvent(m, m.exitListener, l.Presence, presence, "announcePresence") {
209+
return
220210
}
221211
}
222212
}()
@@ -225,15 +215,9 @@ func (m *ListenerManager) announcePresence(presence *PNPresence) {
225215
func (m *ListenerManager) announceFile(file *PNFilesEvent) {
226216
go func() {
227217
lis := m.copyListeners()
228-
229-
AnnounceFileLabel:
230218
for l := range lis {
231-
select {
232-
case <-m.exitListener:
233-
m.pubnub.loggerManager.LogSimple(PNLogLevelTrace, "announceFile: exit listener", false)
234-
break AnnounceFileLabel
235-
236-
case l.File <- file:
219+
if !announceEvent(m, m.exitListener, l.File, file, "announceFile") {
220+
return
237221
}
238222
}
239223
}()

0 commit comments

Comments
 (0)