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
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
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
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
//! The `StreamableHTTP` server-to-client channel (Phase 118.1 plan 10, CONF-07 / G-3).
//!
//! # The property this module exists to preserve
//!
//! `Server::run`'s wiring carries this rationale verbatim
//! (`src/server/mod.rs`), and it is quoted rather than paraphrased because a
//! reader who does not see it will re-introduce the bug:
//!
//! > Transport Actor design (Phase 108, D-01/D-02): the transport is OWNED
//! > by exactly one actor task and is NEVER wrapped in a shared
//! > `Arc<RwLock<T>>`. ALL outbound frames (responses, server-requests,
//! > notifications) funnel through a single UNBOUNDED `send_tx`; inbound
//! > requests go to a SINGLE sequential worker via an UNBOUNDED
//! > `request_tx`. The receive/drain path therefore never blocks on request
//! > execution, request-queue capacity, or a transport write-lock, so an
//! > in-tool `peer.sample()` / `.list_roots()` round-trip cannot deadlock
//! > the loop. Request handling stays serialized (one worker) — zero
//! > behavior change for existing single-request servers.
//!
//! This module reproduces that property over HTTP, where the shape of the hazard
//! is different but the hazard is the same. `dispatch_public_request` holds
//! `state.server.lock().await` for the ENTIRE duration of a tool handler, so a
//! handler awaiting `peer.sample()` holds the server mutex while it waits for
//! the client's answer — and that answer arrives as a POST whose gate
//! (`run_v2_header_gate`, twice) and auth (`extract_and_validate_auth`) stages
//! all take the same mutex. Left alone that is a guaranteed deadlock, and it is
//! also a denial-of-service control: ONE tool call would otherwise take the
//! whole transport offline.
//!
//! Two structural answers, both here:
//!
//! 1. The [`ServerRequestDispatcher`] lives on `ServerState`, OUTSIDE
//! `Mutex<Server>`. A dispatcher behind that mutex is the same deadlock in a
//! different costume.
//! 2. An inbound JSON-RPC RESPONSE is classified and routed BEFORE any stage
//! that takes the mutex (plan 10 Task 3, on both POST entrypoints).
//!
//! # Correlation ownership
//!
//! `Server::run` serves exactly one client, so a server-to-client request has
//! only one place it can go. This transport multiplexes many sessions, so it
//! dispatches through [`ServerRequestDispatcher::dispatch_owned`], which records
//! WHICH session minted each correlation id. The outbound drain resolves that
//! owner to pick the stream, and the inbound path requires it to equal the
//! session presented on the response POST — otherwise a client could resolve
//! another client's pending `sampling/createMessage` (T-118.1-10-02).
//!
//! # Module shape
//!
//! The same three rules `v1_session.rs` states for the v1 pair, for the same
//! reasons: every entry point returns an OWNED answer or performs a whole
//! operation and never hands out a lock guard or an `&Arc<RwLock<..>>`; state
//! lives in a dedicated struct with private fields; any lint allow carries a
//! `// Why:` comment. It is NOT part of the v1/v2 pair — inbound response
//! correlation is era-agnostic — but the OUTBOUND half necessarily routes
//! through `v1::route_to_session_stream`, whose zero-sized twin always answers
//! "no stream", so on a `full-v2` build every outbound dispatch is refused
//! immediately and correctly rather than being silently dropped.
//!
//! # Why this is a submodule and not more of `streamable_http_server.rs`
//!
//! That file is already ~6,700 lines. G-3 does not add to it.
// Why: this is a `pub(crate) mod`, so `pub(crate)` on its items is correct
// (internal-only, never part of the public API) but clippy's nursery
// `redundant_pub_crate` flags it while the crate-level `unreachable_pub` warn
// rejects plain `pub`. The two lints conflict for an internal `pub(crate)`
// module; keeping `pub(crate)` items + this scoped allow is the idiomatic
// resolution already used by `v1_session.rs`, `src/server/task_dispatch.rs` and
// `src/shared/http_body_cap.rs`.
use Arc;
use Duration;
use async_trait;
use StatusCode;
use ;
use Mutex;
use mpsc;
use debug;
use ;
use crate;
use crateListRootsResult;
use crateServerRequestDispatcher;
use cratePeerHandle;
use crateTransportMessage;
use crate;
use crateTransportBackchannel;
use crate;
use crate;
use Deserialize;
/// Capacity of the outbound server-to-client request channel.
///
/// Mirrors `Server::run`'s `mpsc::channel(100)` so the two paths queue alike.
const OUTBOUND_CAPACITY: usize = 100;
/// How long a server-to-client RPC may stay pending on THIS transport.
///
/// EXPLICIT, and deliberately shorter than the 60-second in-process
/// [`DEFAULT_DISPATCH_TIMEOUT`](crate::server::server_request_dispatcher::DEFAULT_DISPATCH_TIMEOUT),
/// because the cost of a pending entry is higher here. Without a bound, a client
/// that opens an SSE stream, triggers N tools that each park on a peer call, and
/// then never answers grows the pending map without limit (T-118.1-10-05).
///
/// And the pending map is not the only thing it pins. `dispatch_public_request`
/// holds `state.server.lock().await` across the whole handler, so every parked
/// peer call also holds the server mutex. A timed-out peer call RELEASES it: the
/// dispatch returns `REQUEST_TIMEOUT`, the handler returns, its stack frame
/// unwinds and `dispatch_public_request`'s guard drops. That is why this value is
/// the transport's real upper bound on how long one absent client can serialize
/// the server — and why it is stated here rather than inherited silently.
///
/// 30 seconds is the compromise: long enough for a real host to run an LLM
/// completion or put an elicitation form in front of a person, short enough that
/// an abandoned round trip is not an outage.
const HTTP_DISPATCH_TIMEOUT: Duration = from_secs;
/// The outbound channel's receiving half, parked until the drain claims it.
///
/// A named alias because the nested generic trips `clippy::type_complexity`
/// inline — and because the shape is load-bearing enough to deserve a name: it is
/// the `(correlation_id, ServerRequest)` pair `Server::run` also queues, held in
/// a `parking_lot::Mutex<Option<..>>` so [`ensure_outbound_drain`] can take it
/// exactly once.
type ParkedOutboundReceiver = ;
/// A one-way notification sink for ONE client session.
///
/// Named because the shape is load-bearing rather than incidental: it is EXACTLY
/// [`ServerProgressReporter::new`](crate::server::progress::ServerProgressReporter::new)'s
/// second parameter, so the transport's sink is handed to the request-scoped
/// progress reporter with no adapter at all. The type is what the two eras share
/// (plan 12 supplies v2's own closure of the same type); the closure is not.
type NotificationSink = ;
/// Logged when an outbound dispatch reaches the drain with no recorded owner.
const NO_OWNER: &str = "outbound dispatch carried no recorded session owner";
/// Logged when the owning session has no live SSE stream to deliver on.
const NO_LIVE_STREAM: &str = "the owning session has no live SSE stream";
// ---------------------------------------------------------------------------
// State.
// ---------------------------------------------------------------------------
/// The transport's server-to-client channel state.
///
/// Lives on `ServerState`, NOT inside `Mutex<Server>` — see the module doc.
/// Both fields are private: every read goes through an operation below, so no
/// caller can take the dispatcher's locks in an order this module did not choose.
pub
/// Hand-written: `ServerRequestDispatcher`'s own `Debug` prints cardinality only,
/// and this must not widen that. It also takes NO lock, for the same reason
/// `V1State`'s does not — a `Debug` impl that blocked inside a panic formatter
/// would turn a diagnostic into a hang.
/// The transport's correlation authority, as an owned handle.
///
/// An OPERATION, not a borrow of the field: the caller gets an `Arc` it owns and
/// this module keeps the only reference to the struct itself.
pub
/// Start the outbound drain, if it is not running already.
///
/// Idempotent: the receiver can only be taken once, so a second call is a no-op.
/// Called from `make_server_state` (covering `pmcp::axum::router()` users) and
/// again from `StreamableHttpServer::start()`, because the first of those is
/// synchronous and may run outside a Tokio runtime — in which case it declines
/// and leaves the receiver for the second.
pub
/// Forward each outbound server-to-client request onto its ORIGINATING session's
/// live SSE stream.
///
/// Exits cleanly when the channel closes — which happens when the last
/// `ServerRequestDispatcher` clone (and with it `outbound_tx`) is dropped, i.e.
/// when the server state goes away. The `Weak` upgrade is the second exit: a
/// message still in flight when the dispatcher dies has nothing left to
/// correlate against.
async
/// Route ONE outbound server-to-client request, or fail its correlation.
///
/// The owner lookup is the step that makes the `(correlation_id, ServerRequest)`
/// channel type sufficient: the session is not carried in the tuple, it is
/// RECORDED against the correlation id at dispatch time and resolved here.
async
// ---------------------------------------------------------------------------
// The session-bound peer handle.
// ---------------------------------------------------------------------------
/// A [`PeerHandle`] bound to ONE v1 session.
///
/// The multiplexing twin of
/// [`DispatchPeerHandle`](crate::server::peer_impl::DispatchPeerHandle): same
/// shared correlation authority, same request construction, same error mapping —
/// the ONLY difference is that every dispatch goes through `dispatch_owned` with
/// this handle's session id, never through `dispatch`. That is what lets the
/// drain send the request to the client that triggered it and lets the inbound
/// path refuse an answer from anyone else.
///
/// Three `Arc` clones per request and no allocation beyond that, matching the
/// cost bar `DispatchPeerHandle`'s rustdoc sets.
pub
/// Hand-written so a `{:?}` of this handle publishes no capability handle and no
/// session token — the same redaction discipline `TransportBackchannel` and
/// `ServerRequestDispatcher` already apply (T-118.1-10-09).
// ---------------------------------------------------------------------------
// Transport-side construction.
// ---------------------------------------------------------------------------
/// Build the V1 one-way notification sink bound to one session's live SSE
/// stream.
///
/// # This closure is v1-ONLY. Only the ATTACHMENT POINT and the TYPE are shared
///
/// [`route_to_session_stream`](v1::route_to_session_stream) is keyed by
/// `session_id`, and v2 has no sessions at all — a v2 GET answers `405`
/// (`sessions_active_truth_table`, `v2_verb_rejection`). Reusing this closure
/// there would look up a session id that cannot exist, `route_to_session_stream`
/// would hand the message straight back on every call, and EVERY v2 progress
/// notification would be silently dropped: a green build with a permanently red
/// `tools-call-with-progress` (T-118.1-11-08).
///
/// So [`attach_session_backchannel`] constructs this only when v1 sessions are
/// live for the request. v2's sink is plan 12's — a bounded per-request queue
/// whose receiver becomes the multi-frame SSE POST response body, which is a
/// different vehicle entirely. What the two eras SHARE is the attachment site
/// and the [`NotificationSink`] type, and that is the whole extent of it.
///
/// # Capture
///
/// Captures `V1State` (which is `Clone` and holds only the SSE-stream map, the
/// session table and the event store) plus the session id — NEVER `ServerState`,
/// which holds `Arc<Mutex<Server>>`. A sink outlives the request that built it,
/// so capturing the server mutex would put a lock handle on the notification
/// path and widen what a leaked closure can reach (T-118.1-11-07).
///
/// Best-effort by construction: a notification is one-way, so a session with no
/// live stream drops it rather than failing a correlation (there is none to
/// fail).
/// Everything [`attach_session_backchannel`] needs about the request it is
/// deciding for.
///
/// A struct rather than four positional arguments because three of the four are
/// booleans-or-options that read identically at a call site.
pub
/// The v1 context this transport would have resolved, had the server opted into
/// v2 era detection.
///
/// Built through the SAME shared unit `fold_v1_handshake_capabilities` uses for
/// its own `None` arm — `first_v1_version` over the server's accept-list — so a
/// context synthesised here is indistinguishable from the one dispatch would
/// have synthesised a moment later. That equality is the whole point: it is what
/// makes attaching a back-channel on a NON-opted-in server a no-op for every
/// other reader of the context.
///
/// Takes the server lock briefly. That is not a new blocking property: the
/// caller sits immediately after `extract_and_validate_auth`, which takes and
/// releases the same lock, and it runs BEFORE dispatch takes it for the handler.
async
/// Attach this request's server-to-client capability handles to its
/// [`ProtocolContext`](crate::types::protocol::ProtocolContext).
///
/// Returns `context` UNCHANGED when there is no back-channel to offer:
///
/// * sessions suppressed for this era — v2 is session-free (HTTP-01), so there
/// is no stream a server-to-client request could be delivered on;
/// * no session on this request — same reason;
/// * the `initialize` handshake itself — the session is being minted by THIS
/// request and no SSE stream can exist yet, so a handle would be inert. Not
/// attaching also keeps `initialize`'s dispatch context exactly as it was.
///
/// # Why an ABSENT context is synthesised rather than skipped
///
/// A stock `Server::builder()` is not opted into v2, so `run_v2_header_gate`
/// short-circuits and this request's `protocol_context` is `None` (D-04's
/// zero-era-code rule) — on precisely the v1 stateful sessions the
/// server-to-client channel exists for. Attaching only to an already-resolved
/// context would therefore give the back-channel to v2-opted-in servers and to
/// nobody else, i.e. to no default deployment at all.
///
/// So a v1 context is synthesised HERE, from the same shared unit that
/// `fold_v1_handshake_capabilities` uses for its own `None` arm. Dispatch already
/// synthesises exactly this value one layer down whenever a v1 handshake has
/// happened, so what reaches the handler is unchanged in every field but the new
/// one — and `era` / `sessions_on` / the legacy-version guard have all already
/// been decided from the ORIGINAL `None` by the time this runs.
///
/// The handles are TRANSPORT-owned and bound to the ORIGINATING session here, at
/// the one site that knows which session this request arrived on. `attach_peer`
/// (plan 11) prefers them over the server's single global `peer_handle`, which is
/// what stops one client's `sampling/createMessage` reaching another's stream
/// (T-118.1-10-04, the T-113-07 class).
///
/// # The v2 branch is INERT here, and that is plan 12's seam
///
/// The `!site.sessions_on` early return means a v2 request gets NO backchannel:
/// no peer (there is no session stream to deliver a server-to-client request on)
/// and no `notification_sink`. A v2 handler therefore finds no reporter and
/// `extra.report_progress(..)` returns `Ok(())` silently — exactly as it did
/// before this phase, so nothing regresses.
///
/// Plan 12 fills that branch at THIS SAME point with a DIFFERENT closure: a
/// bounded per-request queue created before dispatch, whose receiver becomes the
/// multi-frame SSE POST response body. It must not reuse the v1 closure — see
/// [`session_notification_sink`] for why a session-keyed sink silently drops
/// every frame on an era that has no sessions (T-118.1-11-08).
pub async
// ---------------------------------------------------------------------------
// The inbound half: correlation BEFORE the server mutex.
//
// This is the deadlock fix (T-118.1-10-01), not an optimization. See
// `try_route_inbound_response` for the three lock sites it bypasses and for why
// bypassing them is sound for a response and for nothing else.
// ---------------------------------------------------------------------------
/// Try to answer this POST as an inbound JSON-RPC RESPONSE, before the pipeline
/// takes the server mutex.
///
/// Returns `Some(202 Accepted)` for EVERY response envelope — resolved, unknown
/// or wrongly-presented alike — and `None` for every other ingress, which then
/// falls through to the untouched gate / session / auth / dispatch pipeline.
///
/// # Why this must run before `resolve_v2_gate` and `extract_and_validate_auth`
///
/// `dispatch_public_request` holds `state.server.lock().await` for the ENTIRE
/// duration of a tool handler. A handler parked on `peer.sample()` therefore
/// holds the server mutex while it waits for the client's answer — and that
/// answer arrives as a POST whose first stages all take the same mutex:
///
/// * `run_v2_header_gate`, the accept-list read;
/// * `run_v2_header_gate` again, the negotiation read;
/// * `extract_and_validate_auth`, the auth-provider read.
///
/// Left in that order the answer can never reach the dispatcher that would
/// release the handler: a guaranteed deadlock, and a denial-of-service control,
/// because ONE tool call would take the whole transport offline.
///
/// # Why ONLY responses may bypass, and why that is not an auth bypass
///
/// An inbound response carries NO AUTHORITY. It invokes no method, reads nothing
/// and changes no server state: it can only resolve a correlation the SERVER
/// itself minted and is already waiting on. Requests and notifications keep
/// going through the full gate and auth pipeline, unchanged. Do NOT widen this —
/// a request that skipped `extract_and_validate_auth` would be a genuine
/// elevation of privilege (T-118.1-10-06).
///
/// # No new body read
///
/// `HttpIngress::Public(TransportMessage::Response(_))` is the ALREADY-PARSED
/// message, produced by the classifier from the buffer it already owns. The
/// existing `MAX_*` body bounds therefore still apply, unchanged and unmodified
/// (T-118.1-10-07), and this path allocates nothing per unknown id.
/// `pub(super)`, not `pub(crate)`: [`HttpIngress`] is private to the transport
/// module, and a `pub(crate)` signature naming it would be more visible than the
/// type it takes. Both call sites live in that module, so this is exactly wide
/// enough.
pub async
/// Correlate ONE inbound response envelope and answer `202 Accepted`.
///
/// Split out from [`try_route_inbound_response`] so the residual
/// `TransportMessage::Response` arms of the two dispatchers can route through the
/// identical path instead of silently discarding, which is what they used to do.
/// Those arms are unreachable by construction — the classification above answered
/// every response before either dispatcher ran — but a match must stay
/// exhaustive, and an exhaustive arm that DISCARDS is exactly the hole this plan
/// closed.
///
/// # The rejection shape is ONE shape, deliberately
///
/// Three negative cases — an unknown correlation id, an id owned by a DIFFERENT
/// session, and a response presented with no session at all — all answer the
/// same `202 Accepted`, with no body and no distinguishing header, and none of
/// them calls `handle_response`. Differentiating them (say `404` for unknown and
/// `403` for wrong-session) would turn this endpoint into an enumeration oracle
/// for live correlation ids (T-118.1-10-03). The correlation id is logged at
/// `debug`; the payload never is.
pub async