Commit 2f2fa3f
authored
feat(server-subscriptions): expose flatMapMerge concurrency for websocket subscriptions (#2175)
### 📝 Description
`GraphQLWebSocketServer.handleSubscription` pipes the inbound
client-message
flow through `flatMapMerge { ... }` without an explicit `concurrency`
argument,
so it falls back to the kotlinx `DEFAULT_CONCURRENCY = 16`. A single
websocket
session holding more than 16 in-flight subscriptions silently
back-pressures
every subsequent inbound message on the underlying flow — including
`ping`,
`complete`, and additional `subscribe` messages — until one of the 16
in-flight
messages completes. The 17th-and-later subscribes therefore look hung
from the
client's perspective even though the transport is healthy.
This change exposes the concurrency as a configurable value, threaded
through
the three construction paths and both server configurations:
- `GraphQLWebSocketServer` — adds a final constructor parameter
`subscriptionConcurrency: Int = DEFAULT_WS_SUBSCRIPTION_CONCURRENCY`
(top-level `const val = 16`), used as the `concurrency` argument to
`flatMapMerge`.
- `KtorGraphQLWebSocketServer` — forwards a matching parameter to the
superclass constructor.
- `KtorSubscriptionConfiguration` — reads
`graphql.server.subscription.concurrency` from `ApplicationConfig` with
the
same default; `GraphQL.kt` wires it into the Ktor handler.
- `SubscriptionWebSocketHandler` (Spring) — forwards a matching
parameter to
the superclass constructor.
- `SubscriptionConfigurationProperties` — adds
`subscriptionConcurrency: Int = DEFAULT_WS_SUBSCRIPTION_CONCURRENCY` as
a
trailing, defaulted field (preserves data-class binary compatibility);
`SubscriptionGraphQLWsAutoConfiguration` wires it through.
Defaults are unchanged (16), so existing callers see identical
behaviour. Users
who hold many simultaneous subscriptions per session can now raise the
value
(e.g. `Int.MAX_VALUE`) to avoid the back-pressure hang described in the
issue.
A new regression test
(`verify subscription flow honors configured concurrency`) constructs
the
in-memory subscription server with `subscriptionConcurrency = 1` and
sends two
back-to-back `subscribe` messages. Under the previous implicit default
the two
subscriptions would interleave; with `concurrency = 1` the assertion is
that
all four responses for the first subscription id (3 × `next` +
`complete`)
arrive before any response for the second, which is observable and would
have
been impossible without exposing the knob.
The scope is intentionally narrow: only the plumbing and default are
changed.
No change to `TOO_MANY_REQUESTS` handling, no change to graphql-ws
protocol
semantics, no new public types beyond the constant and the trailing
parameters.
### 🔗 Related Issues
Closes #20181 parent f9e0ea4 commit 2f2fa3f
9 files changed
Lines changed: 110 additions & 12 deletions
File tree
- servers
- graphql-kotlin-ktor-server/src/main/kotlin/com/expediagroup/graphql/server/ktor
- subscriptions
- graphql-kotlin-server/src
- main/kotlin/com/expediagroup/graphql/server/execution/subscription
- test/kotlin/com/expediagroup/graphql/server/execution/subscription
- graphql-kotlin-spring-server/src/main/kotlin/com/expediagroup/graphql/server/spring
- subscriptions
servers/graphql-kotlin-ktor-server/src/main/kotlin/com/expediagroup/graphql/server/ktor/GraphQL.kt
Lines changed: 2 additions & 1 deletion
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
169 | 169 | | |
170 | 170 | | |
171 | 171 | | |
172 | | - | |
| 172 | + | |
| 173 | + | |
173 | 174 | | |
174 | 175 | | |
175 | 176 | | |
| |||
Lines changed: 9 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
26 | 26 | | |
27 | 27 | | |
28 | 28 | | |
| 29 | + | |
29 | 30 | | |
30 | 31 | | |
31 | 32 | | |
| |||
284 | 285 | | |
285 | 286 | | |
286 | 287 | | |
| 288 | + | |
| 289 | + | |
| 290 | + | |
| 291 | + | |
| 292 | + | |
| 293 | + | |
| 294 | + | |
| 295 | + | |
287 | 296 | | |
288 | 297 | | |
289 | 298 | | |
| |||
Lines changed: 4 additions & 2 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
17 | 17 | | |
18 | 18 | | |
19 | 19 | | |
| 20 | + | |
20 | 21 | | |
21 | 22 | | |
22 | 23 | | |
| |||
36 | 37 | | |
37 | 38 | | |
38 | 39 | | |
39 | | - | |
| 40 | + | |
| 41 | + | |
40 | 42 | | |
41 | | - | |
| 43 | + | |
42 | 44 | | |
43 | 45 | | |
44 | 46 | | |
| |||
Lines changed: 18 additions & 3 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
57 | 57 | | |
58 | 58 | | |
59 | 59 | | |
60 | | - | |
| 60 | + | |
| 61 | + | |
| 62 | + | |
| 63 | + | |
| 64 | + | |
| 65 | + | |
| 66 | + | |
| 67 | + | |
| 68 | + | |
| 69 | + | |
| 70 | + | |
| 71 | + | |
| 72 | + | |
| 73 | + | |
| 74 | + | |
61 | 75 | | |
62 | 76 | | |
63 | 77 | | |
| |||
67 | 81 | | |
68 | 82 | | |
69 | 83 | | |
70 | | - | |
| 84 | + | |
| 85 | + | |
71 | 86 | | |
72 | 87 | | |
73 | 88 | | |
| |||
86 | 101 | | |
87 | 102 | | |
88 | 103 | | |
89 | | - | |
| 104 | + | |
90 | 105 | | |
91 | 106 | | |
92 | 107 | | |
| |||
Lines changed: 54 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
300 | 300 | | |
301 | 301 | | |
302 | 302 | | |
| 303 | + | |
| 304 | + | |
| 305 | + | |
| 306 | + | |
| 307 | + | |
| 308 | + | |
| 309 | + | |
| 310 | + | |
| 311 | + | |
| 312 | + | |
| 313 | + | |
| 314 | + | |
| 315 | + | |
| 316 | + | |
| 317 | + | |
| 318 | + | |
| 319 | + | |
| 320 | + | |
| 321 | + | |
| 322 | + | |
| 323 | + | |
| 324 | + | |
| 325 | + | |
| 326 | + | |
| 327 | + | |
| 328 | + | |
| 329 | + | |
| 330 | + | |
| 331 | + | |
| 332 | + | |
| 333 | + | |
| 334 | + | |
| 335 | + | |
| 336 | + | |
| 337 | + | |
| 338 | + | |
| 339 | + | |
| 340 | + | |
| 341 | + | |
| 342 | + | |
| 343 | + | |
| 344 | + | |
| 345 | + | |
| 346 | + | |
| 347 | + | |
| 348 | + | |
| 349 | + | |
| 350 | + | |
| 351 | + | |
| 352 | + | |
| 353 | + | |
| 354 | + | |
| 355 | + | |
| 356 | + | |
303 | 357 | | |
304 | 358 | | |
305 | 359 | | |
| |||
Lines changed: 8 additions & 2 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
30 | 30 | | |
31 | 31 | | |
32 | 32 | | |
33 | | - | |
| 33 | + | |
| 34 | + | |
34 | 35 | | |
35 | | - | |
| 36 | + | |
| 37 | + | |
| 38 | + | |
| 39 | + | |
| 40 | + | |
| 41 | + | |
36 | 42 | | |
37 | 43 | | |
38 | 44 | | |
| |||
Lines changed: 9 additions & 1 deletion
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
16 | 16 | | |
17 | 17 | | |
18 | 18 | | |
| 19 | + | |
19 | 20 | | |
20 | 21 | | |
21 | 22 | | |
| |||
91 | 92 | | |
92 | 93 | | |
93 | 94 | | |
94 | | - | |
| 95 | + | |
| 96 | + | |
| 97 | + | |
| 98 | + | |
| 99 | + | |
| 100 | + | |
| 101 | + | |
| 102 | + | |
95 | 103 | | |
96 | 104 | | |
97 | 105 | | |
| |||
Lines changed: 2 additions & 1 deletion
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
63 | 63 | | |
64 | 64 | | |
65 | 65 | | |
66 | | - | |
| 66 | + | |
| 67 | + | |
67 | 68 | | |
68 | 69 | | |
Lines changed: 4 additions & 2 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
17 | 17 | | |
18 | 18 | | |
19 | 19 | | |
| 20 | + | |
20 | 21 | | |
21 | 22 | | |
22 | 23 | | |
| |||
40 | 41 | | |
41 | 42 | | |
42 | 43 | | |
43 | | - | |
| 44 | + | |
| 45 | + | |
44 | 46 | | |
45 | | - | |
| 47 | + | |
46 | 48 | | |
47 | 49 | | |
48 | 50 | | |
| |||
0 commit comments