Skip to content

Commit edf7f1d

Browse files
fix(service): release per-stream subscription connection on unsubscribe
Tie each per-stream subscription's dedicated Postgres connection to its SubscriptionId via an onRemove cleanup Task, released best-effort on removeSubscription so unsubscribe stays total. Fixes the connection leak that exhausted the B1ms connection ceiling under stream-subscription churn. Implements ADR-0063. Closes #683 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent 6423628 commit edf7f1d

4 files changed

Lines changed: 407 additions & 10 deletions

File tree

core/service/Service/EventStore/Postgres/Internal.hs

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -809,9 +809,15 @@ subscribeToStreamEventsImpl ops cfg store entityName streamId callback =
809809
Sessions.selectMaxGlobalPosition
810810
|> Sessions.run conn
811811
|> Task.mapError (toText .> StorageFailure)
812+
813+
-- Release the dedicated per-stream connection when this subscription is
814+
-- removed. Owned by the SubscriptionId, run best-effort on unsubscribe
815+
-- (ADR-0063 §3).
816+
let releaseConnection = Hasql.release connection |> Task.fromIO
817+
812818
subscriptionId <-
813819
store
814-
|> SubscriptionStore.addStreamSubscriptionFromPosition entityName streamId currentMaxPosition callback
820+
|> SubscriptionStore.addStreamSubscriptionWithCleanup entityName streamId currentMaxPosition releaseConnection callback
815821
|> Task.mapError (\err -> SubscriptionError (SubscriptionId "stream") (err |> toText))
816822
Log.debug [fmt|Subscription created: #{toText subscriptionId}|] |> Task.ignoreError
817823
Task.yield subscriptionId

core/service/Service/EventStore/Postgres/SubscriptionStore.hs

Lines changed: 86 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ module Service.EventStore.Postgres.SubscriptionStore (
66
addGlobalSubscriptionFromPosition,
77
addStreamSubscription,
88
addStreamSubscriptionFromPosition,
9+
addStreamSubscriptionWithCleanup,
910
addEntitySubscription,
1011
addEntitySubscriptionFromPosition,
1112
getStreamSubscriptions,
@@ -21,6 +22,7 @@ import Core
2122
import Json qualified
2223
import Log qualified
2324
import Map qualified
25+
import Prelude qualified
2426
import Service.Event (EntityName, Event (..), StreamPosition)
2527
import Service.Event.EventMetadata (EventMetadata (..))
2628
import Service.EventStore.Core (SubscriptionId (..))
@@ -41,9 +43,23 @@ type SubscriptionCallback =
4143
data SubscriptionInfo = SubscriptionInfo
4244
{ callback :: SubscriptionCallback,
4345
startingGlobalPosition :: Maybe StreamPosition,
44-
entityNameFilter :: Maybe EntityName
46+
entityNameFilter :: Maybe EntityName,
47+
-- | Run once when this subscription is removed. Releases any dedicated
48+
-- resource the subscription owns (the per-stream LISTEN connection). A
49+
-- no-op ('Task.yield unit') for subscriptions that own no dedicated resource.
50+
onRemove :: Task Text Unit
4551
}
46-
deriving (Show)
52+
53+
54+
-- Hand-written Show: the callback and onRemove are Task/closure fields with no
55+
-- Show instance, so they are elided. (ADR-0063 §1.)
56+
instance Show SubscriptionInfo where
57+
show info =
58+
"SubscriptionInfo {startingGlobalPosition = "
59+
++ Prelude.show info.startingGlobalPosition
60+
++ ", entityNameFilter = "
61+
++ Prelude.show info.entityNameFilter
62+
++ ", callback = <function>, onRemove = <task>}"
4763

4864

4965
type Subscriptions =
@@ -68,7 +84,7 @@ new = do
6884
addGlobalSubscription :: SubscriptionCallback -> SubscriptionStore -> Task Error SubscriptionId
6985
addGlobalSubscription callback store = do
7086
subId <- Uuid.generate |> Task.map (\result -> result |> toText |> SubscriptionId)
71-
let subscriptionInfo = SubscriptionInfo {callback, startingGlobalPosition = Nothing, entityNameFilter = Nothing}
87+
let subscriptionInfo = SubscriptionInfo {callback, startingGlobalPosition = Nothing, entityNameFilter = Nothing, onRemove = Task.yield unit}
7288
store.globalSubscriptions
7389
|> ConcurrentVar.modify (Map.set subId subscriptionInfo)
7490
Task.yield subId
@@ -78,7 +94,7 @@ addGlobalSubscriptionFromPosition ::
7894
Maybe StreamPosition -> SubscriptionCallback -> SubscriptionStore -> Task Error SubscriptionId
7995
addGlobalSubscriptionFromPosition startingPosition callback store = do
8096
subId <- Uuid.generate |> Task.map (\result -> result |> toText |> SubscriptionId)
81-
let subscriptionInfo = SubscriptionInfo {callback, startingGlobalPosition = startingPosition, entityNameFilter = Nothing}
97+
let subscriptionInfo = SubscriptionInfo {callback, startingGlobalPosition = startingPosition, entityNameFilter = Nothing, onRemove = Task.yield unit}
8298
store.globalSubscriptions
8399
|> ConcurrentVar.modify (Map.set subId subscriptionInfo)
84100
Task.yield subId
@@ -88,7 +104,7 @@ addStreamSubscription ::
88104
EntityName -> StreamId -> SubscriptionCallback -> SubscriptionStore -> Task Error SubscriptionId
89105
addStreamSubscription entityName streamId callback store = do
90106
subId <- Uuid.generate |> Task.map (toText .> SubscriptionId)
91-
let subscriptionInfo = SubscriptionInfo {callback, startingGlobalPosition = Nothing, entityNameFilter = Just entityName}
107+
let subscriptionInfo = SubscriptionInfo {callback, startingGlobalPosition = Nothing, entityNameFilter = Just entityName, onRemove = Task.yield unit}
92108
store.streamSubscriptions |> ConcurrentVar.modify \subscriptionsMap -> do
93109
let currentSubscriptions = subscriptionsMap |> Map.getOrElse streamId Map.empty
94110
subscriptionsMap |> Map.set streamId (currentSubscriptions |> Map.set subId subscriptionInfo)
@@ -104,7 +120,30 @@ addStreamSubscriptionFromPosition ::
104120
Task Error SubscriptionId
105121
addStreamSubscriptionFromPosition entityName streamId startingPosition callback store = do
106122
subId <- Uuid.generate |> Task.map (toText .> SubscriptionId)
107-
let subscriptionInfo = SubscriptionInfo {callback, startingGlobalPosition = startingPosition, entityNameFilter = Just entityName}
123+
let subscriptionInfo = SubscriptionInfo {callback, startingGlobalPosition = startingPosition, entityNameFilter = Just entityName, onRemove = Task.yield unit}
124+
store.streamSubscriptions |> ConcurrentVar.modify \subscriptionsMap -> do
125+
let currentSubscriptions = subscriptionsMap |> Map.getOrElse streamId Map.empty
126+
subscriptionsMap |> Map.set streamId (currentSubscriptions |> Map.set subId subscriptionInfo)
127+
Task.yield subId
128+
129+
130+
addStreamSubscriptionWithCleanup ::
131+
EntityName ->
132+
StreamId ->
133+
Maybe StreamPosition ->
134+
Task Text Unit ->
135+
SubscriptionCallback ->
136+
SubscriptionStore ->
137+
Task Error SubscriptionId
138+
addStreamSubscriptionWithCleanup entityName streamId startingPosition onRemove callback store = do
139+
subId <- Uuid.generate |> Task.map (toText .> SubscriptionId)
140+
let subscriptionInfo =
141+
SubscriptionInfo
142+
{ callback,
143+
startingGlobalPosition = startingPosition,
144+
entityNameFilter = Just entityName,
145+
onRemove
146+
}
108147
store.streamSubscriptions |> ConcurrentVar.modify \subscriptionsMap -> do
109148
let currentSubscriptions = subscriptionsMap |> Map.getOrElse streamId Map.empty
110149
subscriptionsMap |> Map.set streamId (currentSubscriptions |> Map.set subId subscriptionInfo)
@@ -115,7 +154,7 @@ addEntitySubscription ::
115154
EntityName -> SubscriptionCallback -> SubscriptionStore -> Task Error SubscriptionId
116155
addEntitySubscription entityName callback store = do
117156
subId <- Uuid.generate |> Task.map (toText .> SubscriptionId)
118-
let subscriptionInfo = SubscriptionInfo {callback, startingGlobalPosition = Nothing, entityNameFilter = Just entityName}
157+
let subscriptionInfo = SubscriptionInfo {callback, startingGlobalPosition = Nothing, entityNameFilter = Just entityName, onRemove = Task.yield unit}
119158
store.entitySubscriptions |> ConcurrentVar.modify \subscriptionsMap -> do
120159
let currentSubscriptions = subscriptionsMap |> Map.getOrElse entityName Map.empty
121160
subscriptionsMap |> Map.set entityName (currentSubscriptions |> Map.set subId subscriptionInfo)
@@ -130,7 +169,7 @@ addEntitySubscriptionFromPosition ::
130169
Task Error SubscriptionId
131170
addEntitySubscriptionFromPosition entityName startingPosition callback store = do
132171
subId <- Uuid.generate |> Task.map (toText .> SubscriptionId)
133-
let subscriptionInfo = SubscriptionInfo {callback, startingGlobalPosition = startingPosition, entityNameFilter = Just entityName}
172+
let subscriptionInfo = SubscriptionInfo {callback, startingGlobalPosition = startingPosition, entityNameFilter = Just entityName, onRemove = Task.yield unit}
134173
store.entitySubscriptions |> ConcurrentVar.modify \subscriptionsMap -> do
135174
let currentSubscriptions = subscriptionsMap |> Map.getOrElse entityName Map.empty
136175
subscriptionsMap |> Map.set entityName (currentSubscriptions |> Map.set subId subscriptionInfo)
@@ -209,8 +248,39 @@ dispatch streamId message store = do
209248
allCallbacks |> AsyncTask.runAllIgnoringErrors
210249

211250

251+
-- | Find a subscription's 'onRemove' cleanup action across the global, stream,
252+
-- and entity maps. Returns 'Task.yield unit' when the id is unknown, so removing
253+
-- an already-removed or never-registered id is a no-op rather than an error.
254+
findOnRemove :: SubscriptionId -> SubscriptionStore -> Task Error (Task Text Unit)
255+
findOnRemove subId store = do
256+
globalSubs <- store.globalSubscriptions |> ConcurrentVar.peek
257+
streamSubsMap <- store.streamSubscriptions |> ConcurrentVar.peek
258+
entitySubsMap <- store.entitySubscriptions |> ConcurrentVar.peek
259+
-- Search the global map, then every stream's map, then every entity's map for
260+
-- the id. SubscriptionIds are unique, so at most one map holds it.
261+
let lookupIn subs acc =
262+
case acc of
263+
Just info -> Just info
264+
Nothing -> subs |> Map.get subId
265+
let fromStream = streamSubsMap |> Map.values |> Array.reduce lookupIn Nothing
266+
let fromEntity = entitySubsMap |> Map.values |> Array.reduce lookupIn Nothing
267+
let found =
268+
case globalSubs |> Map.get subId of
269+
Just info -> Just info
270+
Nothing -> case fromStream of
271+
Just info -> Just info
272+
Nothing -> fromEntity
273+
case found of
274+
Just info -> Task.yield info.onRemove
275+
Nothing -> Task.yield (Task.yield unit)
276+
277+
212278
removeSubscription :: SubscriptionId -> SubscriptionStore -> Task Error Unit
213279
removeSubscription subId store = do
280+
-- Read the subscription's cleanup BEFORE deleting, so the cleanup is taken
281+
-- (not left in place) and a double-remove cannot run it twice.
282+
cleanup <- store |> findOnRemove subId
283+
214284
-- Remove from global subscriptions
215285
store.globalSubscriptions
216286
|> ConcurrentVar.modify (Map.remove subId)
@@ -225,4 +295,11 @@ removeSubscription subId store = do
225295
subscriptionsMap
226296
|> Map.mapValues (\entitySubs -> entitySubs |> Map.remove subId)
227297

228-
Task.yield ()
298+
-- Run the cleanup best-effort: a failed release is logged and swallowed so
299+
-- unsubscribe stays total (ADR-0063 §2, design goal 4).
300+
cleanup
301+
|> Task.recover \err -> do
302+
Log.warn [fmt|Subscription #{toText subId} cleanup failed: #{err}|]
303+
|> Task.ignoreError
304+
Task.yield unit
305+
|> Task.mapError (\_ -> OtherError)

0 commit comments

Comments
 (0)