asupersync 0.4.10

Spec-first, cancel-correct, capability-secure async runtime for Rust.
Documentation
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
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
//! High-level QUIC connection API (`arq-quic-epic-b0k8qo.1.6`, "A6").
//!
//! The lower `quic_native` modules expose the raw data-plane machinery: the
//! [`NativeQuicConnection`] state machine ([`connection`](super::connection)),
//! the stream reassembly + flow-control table ([`streams`](super::streams)), the
//! 1-RTT packet-protection boundary
//! ([`crate::net::atp::quic::packet_protection`]), the real UDP endpoint
//! ([`endpoint`](super::endpoint)), and the integrated event loop
//! ([`managed_endpoint`](super::managed_endpoint)). What was missing was a
//! *usable application surface* on top of those pieces — the "thin adapter
//! target" that Phase B (`transport_quic`) needs so it can issue ~5-7
//! high-level calls instead of hand-driving packet/crypto/ACK loops.
//!
//! This module provides that surface as [`QuicConnection`]: a role-aware handle
//! that drives the handshake, sends/receives RFC 9221 DATAGRAMs (A1/A2), opens
//! a bidirectional control stream as ordered reliable bytes (A3), exposes path
//! statistics (RTT / cwnd / loss) for the Phase C controller, and closes
//! gracefully.
//!
//! # Transport model (no-claim boundary)
//!
//! [`QuicConnection`] is the application-facing handle used by both production
//! and deterministic transports. What is deliberately *deterministic /
//! lab-only* is the in-memory transport that carries bytes between two handles
//! — [`establish_loopback`] and [`pump_app_data`]. Production callers instead
//! bind the same handle to [`NativeQuicUdpConnection`](super::udp_connection::NativeQuicUdpConnection),
//! which owns:
//!
//! * a real UDP endpoint and caller-driven bounded I/O/timer operations;
//! * the actual local and authenticated peer connection IDs;
//! * the completed rustls QUIC handshake, negotiated ALPN, and handshake-derived
//!   1-RTT packet protection.
//!
//! The loopback helpers still drive handshake transitions directly; they are an
//! honest deterministic transport, not a wire-security proof. Their DATAGRAM
//! and STREAM bytes flow through
//! [`NativeQuicConnection::generate_frames`] →
//! [`NativeQuicConnection::process_packet_payload`]. The separate live-UDP
//! integration proof interposes the real handshake and protect/unprotect path.
//! Neither surface claims a multi-connection listener, external QUIC/H3
//! interoperability, migration, 0-RTT, or production deployment readiness.
//!
//! # Fail-closed
//!
//! The client identity gate is preserved end to end: a client
//! [`QuicConnection`] cannot reach [`QuicConnectionState::Established`] unless
//! the application has recorded a verified server identity via
//! [`QuicConnection::record_verified_server_identity`] (which, on the production
//! path, must follow a genuine in-handshake X.509 verification — there is no
//! insecure skip-verify default, per `asupersync-7pwwwe` / `b0k8qo.1.5`).
//! [`establish_loopback`] does not bypass this: it propagates the
//! [`QuicTlsError::ServerCertificateUnverified`](super::tls::QuicTlsError)
//! fail-closed error if the client identity was not recorded first.

use crate::bytes::{Bytes, BytesMut};
use crate::cx::Cx;
use crate::net::atp::protocol::quic_frames::QuicFrame;
use std::task::{Context as TaskContext, Poll};

use super::connection::{
    NativeQuicConnection, NativeQuicConnectionConfig, NativeQuicConnectionError,
};
use super::streams::{StreamId, StreamReadiness, StreamRole};
use super::transport::{PacketNumberSpace, QuicConnectionState};

/// Opt-in stderr tracing for the high-level QUIC API, gated by `ATP_QUIC_TRACE`
/// so the production path stays silent (mirrors `connection.rs`'s `quictrace!`).
macro_rules! apitrace {
    ($($arg:tt)*) => {
        if std::env::var_os("ATP_QUIC_TRACE").is_some() {
            eprintln!("[atp-quic-api] {}", format!($($arg)*));
        }
    };
}

/// Default maximum 1-RTT packet payload budget used by the deterministic
/// loopback transport when a caller does not specify one.
///
/// A datagram frame is bounded to 1200 bytes by the connection, so this
/// comfortably carries a full datagram plus framing while still exercising the
/// multi-packet path for large stream transfers.
pub const DEFAULT_MAX_PACKET_BYTES: usize = 1350;

/// Safety cap on [`pump_until_idle`] iterations so a misbehaving queue can never
/// spin forever; far above any realistic per-direction flight.
const PUMP_ITERATION_CAP: usize = 16_384;

/// Point-in-time path statistics for the Phase C adaptive controller.
///
/// All values are read from the connection's loss-recovery / congestion-control
/// state machine; they are advisory signals, not guarantees.
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct QuicPathStats {
    /// Smoothed RTT estimate in microseconds, once at least one RTT sample has
    /// been observed.
    pub smoothed_rtt_micros: Option<u64>,
    /// Most recent RTT sample in microseconds.
    pub latest_rtt_micros: Option<u64>,
    /// RTT variation in microseconds.
    pub rttvar_micros: Option<u64>,
    /// Current congestion window in bytes.
    pub congestion_window_bytes: u64,
    /// Bytes currently in flight (sent, not yet acknowledged or declared lost).
    pub bytes_in_flight: u64,
    /// Probe-timeout backoff count (a rough loss / tail-latency signal).
    pub pto_count: u32,
    /// Cumulative packets acknowledged by the recovery state.
    pub packets_acked: u64,
    /// Cumulative packets declared lost by the recovery state.
    pub packets_lost: u64,
    /// Cumulative packet loss rate over packets that reached an acked/lost
    /// recovery outcome.
    pub loss_rate: f64,
}

/// A high-level, role-aware handle over a single native QUIC connection.
///
/// This is the application-facing surface for the QUIC data plane: drive the
/// handshake, then send/receive datagrams and reliable control-stream bytes
/// without touching frames, packets, or crypto directly. See the [module
/// docs](self) for the transport model and fail-closed semantics.
#[derive(Debug)]
pub struct QuicConnection {
    inner: NativeQuicConnection,
    role: StreamRole,
    /// Monotonic 1-RTT packet number used by the deterministic loopback
    /// transport. The production path assigns packet numbers in the protect
    /// layer instead.
    next_app_pn: u64,
}

impl QuicConnection {
    fn from_config(mut config: NativeQuicConnectionConfig, role: StreamRole) -> Self {
        config.role = role;
        Self {
            inner: NativeQuicConnection::new(config),
            role,
            next_app_pn: 0,
        }
    }

    /// Construct a client-role connection handle.
    ///
    /// The `role` field of `config` is forced to [`StreamRole::Client`].
    #[must_use]
    pub fn client(config: NativeQuicConnectionConfig) -> Self {
        apitrace!("event=conn_new role=client");
        Self::from_config(config, StreamRole::Client)
    }

    /// Construct a server-role connection handle.
    ///
    /// The `role` field of `config` is forced to [`StreamRole::Server`].
    #[must_use]
    pub fn server(config: NativeQuicConnectionConfig) -> Self {
        apitrace!("event=conn_new role=server");
        Self::from_config(config, StreamRole::Server)
    }

    /// This connection's role.
    #[must_use]
    pub fn role(&self) -> StreamRole {
        self.role
    }

    /// Current transport state.
    #[must_use]
    pub fn state(&self) -> QuicConnectionState {
        self.inner.state()
    }

    /// Whether application (1-RTT) data may be sent right now (handshake is
    /// confirmed and 1-RTT keys are installed).
    #[must_use]
    pub fn can_send_app_data(&self) -> bool {
        self.inner.can_send_1rtt()
    }

    /// Borrow the underlying state machine for advanced/diagnostic use.
    #[must_use]
    pub fn inner(&self) -> &NativeQuicConnection {
        &self.inner
    }

    /// Mutably borrow the state machine for the crate's wire drivers.
    ///
    /// This stays crate-private so application code cannot bypass the
    /// role-aware stream API or replay handshake transitions manually.
    #[cfg(feature = "tls")]
    pub(crate) fn inner_mut(&mut self) -> &mut NativeQuicConnection {
        &mut self.inner
    }

    // -- handshake -----------------------------------------------------------

    /// Begin the handshake (`Idle → Handshaking`).
    ///
    /// # Errors
    /// Returns a [`NativeQuicConnectionError`] if `cx` is cancelled or the
    /// transport rejects the transition from its current state.
    pub fn begin_handshake(&mut self, cx: &Cx) -> Result<(), NativeQuicConnectionError> {
        apitrace!("event=begin_handshake role={:?}", self.role);
        self.inner.begin_handshake(cx)
    }

    /// Mark handshake-level keys installed.
    ///
    /// # Errors
    /// Returns a [`NativeQuicConnectionError`] if `cx` is cancelled or the TLS
    /// machine rejects installing handshake keys in its current level.
    pub fn mark_handshake_keys_available(
        &mut self,
        cx: &Cx,
    ) -> Result<(), NativeQuicConnectionError> {
        self.inner.on_handshake_keys_available(cx)
    }

    /// Mark 1-RTT application keys installed.
    ///
    /// # Errors
    /// Returns a [`NativeQuicConnectionError`] if `cx` is cancelled or the TLS
    /// machine rejects installing 1-RTT keys in its current level.
    pub fn mark_app_keys_available(&mut self, cx: &Cx) -> Result<(), NativeQuicConnectionError> {
        self.inner.on_1rtt_keys_available(cx)
    }

    /// Record that the application has verified the server's identity.
    ///
    /// This is the only way to clear the fail-closed server-identity gate on a
    /// **client** connection. On the production path it must be called only
    /// after a genuine certificate verification (chain + hostname + signature)
    /// has succeeded; there is deliberately no insecure skip-verify toggle
    /// (`asupersync-7pwwwe`). On a server connection it has no effect on the
    /// handshake — the identity gate is only consulted on the client path.
    pub fn record_verified_server_identity(&mut self) {
        apitrace!("event=server_identity_recorded role={:?}", self.role);
        self.inner.record_verified_server_identity();
    }

    /// Confirm the handshake (`Handshaking → Established`).
    ///
    /// # Errors
    /// For a client connection this fails closed with
    /// [`QuicTlsError::ServerCertificateUnverified`](super::tls::QuicTlsError)
    /// (wrapped in [`NativeQuicConnectionError::Tls`]) unless
    /// [`Self::record_verified_server_identity`] was called first. It also
    /// returns a [`NativeQuicConnectionError`] if `cx` is cancelled or the
    /// transport/TLS state is not ready to confirm.
    pub fn confirm_handshake(&mut self, cx: &Cx) -> Result<(), NativeQuicConnectionError> {
        let result = self.inner.on_handshake_confirmed(cx);
        match &result {
            Ok(()) => apitrace!("event=handshake_confirmed role={:?}", self.role),
            Err(err) => apitrace!(
                "event=handshake_confirm_failed role={:?} err={err:?}",
                self.role
            ),
        }
        result
    }

    // -- datagrams (A1 / A2) -------------------------------------------------

    /// Queue an unreliable application datagram (RFC 9221) for transmission.
    ///
    /// The handle enforces that application data may only be sent once the
    /// connection is established (1-RTT). The payload is bounded by the
    /// connection's `max_datagram_frame_size`; an oversize payload is rejected
    /// fail-closed.
    ///
    /// # Errors
    /// * [`NativeQuicConnectionError::InvalidState`] if the connection is not
    ///   yet established.
    /// * [`NativeQuicConnectionError::DatagramTooLarge`] if the encoded frame
    ///   would exceed the maximum datagram frame size.
    /// * [`NativeQuicConnectionError::Cancelled`] if `cx` is cancelled.
    pub fn send_datagram(
        &mut self,
        cx: &Cx,
        payload: Bytes,
    ) -> Result<(), NativeQuicConnectionError> {
        if !self.inner.can_send_1rtt() {
            apitrace!(
                "event=datagram_send_reject reason=not_established role={:?}",
                self.role
            );
            return Err(NativeQuicConnectionError::InvalidState(
                "send_datagram requires an established 1-RTT connection",
            ));
        }
        let len = payload.len();
        let result = self.inner.send_datagram(cx, payload);
        if result.is_ok() {
            apitrace!("event=datagram_queued role={:?} len={len}", self.role);
        }
        result
    }

    /// Pop the next received datagram payload, if any (non-blocking).
    #[must_use]
    pub fn recv_datagram(&mut self) -> Option<Bytes> {
        self.inner.recv_datagram()
    }

    /// Cx-aware poll for the next received datagram payload (no busy-poll).
    ///
    /// Returns [`Poll::Ready`] with the payload when one is buffered, observes
    /// cancellation through `cx` (returning [`Poll::Ready`] with
    /// [`NativeQuicConnectionError::Cancelled`]), and otherwise registers
    /// `task_cx`'s waker for the next datagram arrival and returns
    /// [`Poll::Pending`].
    pub fn poll_recv_datagram(
        &mut self,
        cx: &Cx,
        task_cx: &mut TaskContext<'_>,
    ) -> Poll<Result<Bytes, NativeQuicConnectionError>> {
        self.inner.poll_recv_datagram(cx, task_cx)
    }

    /// Number of received datagram payloads currently buffered.
    #[must_use]
    pub fn pending_datagram_count(&self) -> usize {
        self.inner.pending_datagram_count()
    }

    /// Total datagrams emitted onto the (loopback) wire.
    #[must_use]
    pub fn datagrams_sent(&self) -> u64 {
        self.inner.datagrams_sent()
    }

    /// Total datagrams accepted on receive (counted before any drop-oldest).
    #[must_use]
    pub fn datagrams_received(&self) -> u64 {
        self.inner.datagrams_received()
    }

    // -- control stream (A3) -------------------------------------------------

    /// Open a bidirectional control stream and return its id.
    ///
    /// # Errors
    /// Returns a [`NativeQuicConnectionError`] if `cx` is cancelled, the
    /// connection is not in a data-transfer state, or the local stream limit is
    /// exhausted.
    pub fn open_control_stream(&mut self, cx: &Cx) -> Result<StreamId, NativeQuicConnectionError> {
        let id = self.open_bidi_stream(cx)?;
        apitrace!(
            "event=control_stream_open role={:?} stream={}",
            self.role,
            id.0
        );
        Ok(id)
    }

    /// Open a locally initiated bidirectional application stream.
    pub fn open_bidi_stream(&mut self, cx: &Cx) -> Result<StreamId, NativeQuicConnectionError> {
        self.inner.open_local_bidi(cx)
    }

    /// Open a locally initiated unidirectional application stream.
    pub fn open_uni_stream(&mut self, cx: &Cx) -> Result<StreamId, NativeQuicConnectionError> {
        self.inner.open_local_uni(cx)
    }

    /// Queue reliable, ordered bytes (and an optional FIN) on a control stream.
    ///
    /// # Errors
    /// Returns a [`NativeQuicConnectionError`] if `cx` is cancelled, the
    /// connection is not in a data-transfer state, the stream is unknown, or the
    /// flow-control window is exhausted (a `STREAM_DATA_BLOCKED` is queued).
    pub fn write_control(
        &mut self,
        cx: &Cx,
        stream: StreamId,
        data: Bytes,
        fin: bool,
    ) -> Result<(), NativeQuicConnectionError> {
        self.inner.write_stream_bytes(cx, stream, data, fin)
    }

    /// Queue reliable ordered bytes on any locally writable stream.
    pub fn write_stream(
        &mut self,
        cx: &Cx,
        stream: StreamId,
        data: Bytes,
        fin: bool,
    ) -> Result<(), NativeQuicConnectionError> {
        self.inner.write_stream_bytes(cx, stream, data, fin)
    }

    /// Whether one stream still has bytes or FIN queued for packet assembly.
    #[must_use]
    pub fn has_pending_stream_frames(&self, stream: StreamId) -> bool {
        self.inner.has_pending_stream_frames_for(stream)
    }

    /// Number of application payload bytes queued for one stream.
    #[must_use]
    pub fn pending_stream_data_bytes(&self, stream: StreamId) -> u64 {
        self.inner.pending_stream_data_bytes_for(stream)
    }

    /// Poll the payload capacity available under both stream and connection
    /// flow control.
    pub fn poll_stream_send_capacity(
        &mut self,
        cx: &Cx,
        stream: StreamId,
        required: u64,
        task_cx: &mut TaskContext<'_>,
    ) -> Poll<Result<u64, NativeQuicConnectionError>> {
        self.inner
            .poll_stream_send_capacity(cx, stream, required, task_cx)
    }

    /// Poll until prior STREAM writes have drained and both flow-control
    /// scopes admit `required` bytes for a new application frame.
    pub fn poll_stream_write_ready(
        &mut self,
        cx: &Cx,
        stream: StreamId,
        required: u64,
        task_cx: &mut TaskContext<'_>,
    ) -> Poll<Result<u64, NativeQuicConnectionError>> {
        self.inner
            .poll_stream_write_ready(cx, stream, required, task_cx)
    }

    /// Poll until every queued STREAM frame for one stream leaves packet
    /// assembly, including a zero-byte terminal FIN.
    pub fn poll_stream_queue_drained(
        &mut self,
        cx: &Cx,
        stream: StreamId,
        task_cx: &mut TaskContext<'_>,
    ) -> Poll<Result<(), NativeQuicConnectionError>> {
        self.inner.poll_stream_queue_drained(cx, stream, task_cx)
    }

    /// Read contiguous reassembled bytes from a control stream (up to `max`).
    ///
    /// Returns an empty buffer when no further contiguous bytes are available
    /// yet; check [`Self::is_control_eof`] to distinguish "more later" from EOF.
    ///
    /// # Errors
    /// Returns a [`NativeQuicConnectionError`] if `cx` is cancelled, the
    /// connection is not in a stream-active state, or the stream is unknown.
    pub fn read_control(
        &mut self,
        cx: &Cx,
        stream: StreamId,
        max: usize,
    ) -> Result<Bytes, NativeQuicConnectionError> {
        self.inner.read_stream_bytes(cx, stream, max)
    }

    /// Read contiguous bytes from any readable stream.
    pub fn read_stream(
        &mut self,
        cx: &Cx,
        stream: StreamId,
        max: usize,
    ) -> Result<Bytes, NativeQuicConnectionError> {
        self.inner.read_stream_bytes(cx, stream, max)
    }

    /// Configure a bounded receive window and queue MAX_STREAM_DATA.
    pub fn configure_stream_receive_window(
        &mut self,
        cx: &Cx,
        stream: StreamId,
        window: u64,
    ) -> Result<u64, NativeQuicConnectionError> {
        self.inner.configure_stream_recv_window(cx, stream, window)
    }

    /// Increase the connection receive limit and queue MAX_DATA.
    pub fn advertise_connection_receive_limit(
        &mut self,
        cx: &Cx,
        limit: u64,
    ) -> Result<(), NativeQuicConnectionError> {
        self.inner.advertise_connection_recv_limit(cx, limit)
    }

    /// Clamp one stream's send limit for a deterministic flow-control test.
    #[cfg(any(test, feature = "test-internals"))]
    pub fn constrain_stream_send_limit_for_testing(
        &mut self,
        cx: &Cx,
        stream: StreamId,
        limit: u64,
    ) -> Result<u64, NativeQuicConnectionError> {
        self.inner
            .constrain_stream_send_limit_for_testing(cx, stream, limit)
    }

    /// Clamp connection send credit for a deterministic flow-control test.
    #[cfg(any(test, feature = "test-internals"))]
    pub fn constrain_connection_send_limit_for_testing(
        &mut self,
        cx: &Cx,
        limit: u64,
    ) -> Result<u64, NativeQuicConnectionError> {
        self.inner
            .constrain_connection_send_limit_for_testing(cx, limit)
    }

    /// Consume the next deterministic stream-readiness edge, if one exists.
    pub fn next_readable_stream(
        &mut self,
        cx: &Cx,
    ) -> Result<Option<StreamReadiness>, NativeQuicConnectionError> {
        self.inner.next_readable_stream(cx)
    }

    /// Poll for the next stream-readiness edge without busy-polling.
    pub fn poll_next_readable_stream(
        &mut self,
        cx: &Cx,
        task_cx: &mut TaskContext<'_>,
    ) -> Poll<Result<StreamReadiness, NativeQuicConnectionError>> {
        self.inner.poll_next_readable_stream(cx, task_cx)
    }

    /// Queue a local RESET_STREAM for one stream send side.
    pub fn reset_stream(
        &mut self,
        cx: &Cx,
        stream: StreamId,
        app_error_code: u64,
    ) -> Result<(), NativeQuicConnectionError> {
        let final_size = self
            .inner
            .streams()
            .stream(stream)
            .map_err(NativeQuicConnectionError::from)?
            .send_offset;
        self.inner
            .reset_stream_send(cx, stream, app_error_code, final_size)
    }

    /// Queue a local STOP_SENDING for one stream receive side.
    pub fn stop_stream_receiving(
        &mut self,
        cx: &Cx,
        stream: StreamId,
        app_error_code: u64,
    ) -> Result<(), NativeQuicConnectionError> {
        self.inner.stop_receiving(cx, stream, app_error_code)
    }

    /// Whether the application has consumed a control stream through its FIN.
    ///
    /// # Errors
    /// Returns a [`NativeQuicConnectionError`] if the stream is unknown.
    pub fn is_control_eof(&self, stream: StreamId) -> Result<bool, NativeQuicConnectionError> {
        self.inner.is_stream_read_eof(stream)
    }

    /// Whether the application has consumed a stream through peer FIN.
    pub fn is_stream_eof(&self, stream: StreamId) -> Result<bool, NativeQuicConnectionError> {
        self.inner.is_stream_read_eof(stream)
    }

    /// Bounded stream prefix retained when peer RESET_STREAM discarded
    /// receive data before an application protocol classified the stream.
    pub fn reset_stream_buffered_prefix(
        &self,
        stream: StreamId,
    ) -> Result<Bytes, NativeQuicConnectionError> {
        self.inner.reset_stream_buffered_prefix(stream)
    }

    // -- path stats (Phase C) ------------------------------------------------

    /// Snapshot of path statistics for adaptive control.
    #[must_use]
    pub fn path_stats(&self) -> QuicPathStats {
        let transport = self.inner.transport();
        let rtt = transport.rtt();
        QuicPathStats {
            smoothed_rtt_micros: rtt.smoothed_rtt_micros(),
            latest_rtt_micros: rtt.latest_rtt_micros(),
            rttvar_micros: rtt.rttvar_micros(),
            congestion_window_bytes: transport.congestion_window_bytes(),
            bytes_in_flight: transport.bytes_in_flight(),
            pto_count: transport.pto_count(),
            packets_acked: transport.packets_acked_total(),
            packets_lost: transport.packets_lost_total(),
            loss_rate: transport.packet_loss_rate(),
        }
    }

    // -- close ---------------------------------------------------------------

    /// Begin a graceful close (enter draining with an application error code).
    ///
    /// # Errors
    /// Returns a [`NativeQuicConnectionError`] if `cx` is cancelled or the
    /// transport rejects the transition.
    pub fn begin_close(
        &mut self,
        cx: &Cx,
        now_micros: u64,
        app_error_code: u64,
    ) -> Result<(), NativeQuicConnectionError> {
        apitrace!(
            "event=begin_close role={:?} code={app_error_code}",
            self.role
        );
        self.inner.begin_close(cx, now_micros, app_error_code)
    }
}

/// Drive a client + server [`QuicConnection`] pair through the handshake to
/// [`QuicConnectionState::Established`] using the deterministic in-memory
/// transport.
///
/// This is the lab/test substitute for the production event loop + real
/// wire-CRYPTO handshake driver (see the [module docs](self)). It drives the
/// key-availability transitions directly (matching `tests/quic_h3_e2e.rs`); the
/// real driver will instead advance them from exchanged CRYPTO bytes.
///
/// The fail-closed client-identity gate is preserved: the `client` must have
/// recorded a verified server identity via
/// [`QuicConnection::record_verified_server_identity`] before this call, or it
/// returns the wrapped
/// [`QuicTlsError::ServerCertificateUnverified`](super::tls::QuicTlsError).
///
/// # Errors
/// Propagates any [`NativeQuicConnectionError`] from the underlying transitions,
/// including the fail-closed identity gate.
pub fn establish_loopback(
    cx: &Cx,
    client: &mut QuicConnection,
    server: &mut QuicConnection,
) -> Result<(), NativeQuicConnectionError> {
    client.begin_handshake(cx)?;
    server.begin_handshake(cx)?;
    client.mark_handshake_keys_available(cx)?;
    server.mark_handshake_keys_available(cx)?;
    client.mark_app_keys_available(cx)?;
    server.mark_app_keys_available(cx)?;
    // Server confirms freely; client must have a recorded verified identity or
    // it fails closed here.
    server.confirm_handshake(cx)?;
    client.confirm_handshake(cx)?;
    apitrace!("event=loopback_established");
    Ok(())
}

/// Deterministic in-memory transport: drain one 1-RTT packet's worth of pending
/// application frames (control STREAM + DATAGRAM) from `from` and deliver them
/// to `to`, returning the number of frames moved.
///
/// This stands in for the production UDP send + AEAD protect/unprotect path; the
/// bytes themselves really flow through
/// [`NativeQuicConnection::generate_frames`] →
/// [`NativeQuicConnection::process_packet_payload`].
///
/// It deliberately does not record sent packets through the loss-recovery
/// machine (`on_packet_sent` / `on_ack_received`), so RTT and bytes-in-flight
/// reported by [`QuicConnection::path_stats`] stay at their initial defaults
/// under the loopback; the production event-loop path populates those signals.
///
/// # Errors
/// Returns a [`NativeQuicConnectionError`] if frame generation, encoding, or
/// payload processing fails (including `cx` cancellation).
pub fn pump_app_data(
    cx: &Cx,
    from: &mut QuicConnection,
    to: &mut QuicConnection,
    max_packet_bytes: usize,
    now_micros: u64,
) -> Result<usize, NativeQuicConnectionError> {
    let frames: Vec<QuicFrame> =
        from.inner
            .generate_frames(cx, PacketNumberSpace::ApplicationData, max_packet_bytes)?;
    if frames.is_empty() {
        return Ok(0);
    }
    let mut payload = BytesMut::new();
    for frame in &frames {
        frame.encode(&mut payload)?;
    }
    let packet_number = from.next_app_pn;
    from.next_app_pn = from.next_app_pn.saturating_add(1);
    let delivery = to.inner.process_packet_payload(
        cx,
        PacketNumberSpace::ApplicationData,
        packet_number,
        &payload,
        now_micros,
    );
    match delivery {
        Ok(()) => from.inner.on_generated_frames_delivered(&frames)?,
        Err(error) => {
            from.inner.on_generated_frames_dropped(&frames)?;
            return Err(error);
        }
    }
    apitrace!(
        "event=pump frames={} bytes={} pn={packet_number}",
        frames.len(),
        payload.len()
    );
    Ok(frames.len())
}

/// Deterministically drop one generated application packet and requeue every
/// reliable frame it carried through the same recovery ledger used by the
/// production native connection.
///
/// DATAGRAM frames remain intentionally lost. The helper exists for causal lab
/// and integration tests of STREAM/RESET_STREAM/STOP_SENDING recovery without
/// sockets, sleeps, or probabilistic loss injection.
#[cfg(any(test, feature = "test-internals"))]
pub fn drop_app_data_packet(
    cx: &Cx,
    from: &mut QuicConnection,
    max_packet_bytes: usize,
) -> Result<usize, NativeQuicConnectionError> {
    let frames =
        from.inner
            .generate_frames(cx, PacketNumberSpace::ApplicationData, max_packet_bytes)?;
    if frames.is_empty() {
        return Ok(0);
    }
    from.next_app_pn = from.next_app_pn.saturating_add(1);
    from.inner.on_generated_frames_dropped(&frames)?;
    Ok(frames.len())
}

/// Repeatedly [`pump_app_data`] from `from` to `to` until `from` has no further
/// pending application frames, returning the total number of frames moved.
///
/// Bounded by `PUMP_ITERATION_CAP` so a stuck queue can never spin forever.
///
/// # Errors
/// Returns a [`NativeQuicConnectionError`] on the first failing pump round, or
/// [`NativeQuicConnectionError::InvalidState`] if the iteration cap is hit
/// (which would indicate a non-draining queue).
pub fn pump_until_idle(
    cx: &Cx,
    from: &mut QuicConnection,
    to: &mut QuicConnection,
    max_packet_bytes: usize,
    now_micros: u64,
) -> Result<usize, NativeQuicConnectionError> {
    let mut total = 0;
    for _ in 0..PUMP_ITERATION_CAP {
        let moved = pump_app_data(cx, from, to, max_packet_bytes, now_micros)?;
        if moved == 0 {
            return Ok(total);
        }
        total += moved;
    }
    Err(NativeQuicConnectionError::InvalidState(
        "pump_until_idle exceeded its iteration cap without draining",
    ))
}

#[cfg(test)]
mod tests {
    #![allow(clippy::cast_possible_truncation)]
    use super::*;

    fn test_cx() -> Cx<crate::cx::cap::All> {
        Cx::for_testing()
    }

    fn pair() -> (QuicConnection, QuicConnection) {
        // `NativeQuicConnectionConfig` is `Copy`, so each constructor copies it.
        let cfg = NativeQuicConnectionConfig::default();
        (QuicConnection::client(cfg), QuicConnection::server(cfg))
    }

    fn established_pair(cx: &Cx) -> (QuicConnection, QuicConnection) {
        let (mut client, mut server) = pair();
        client.record_verified_server_identity();
        establish_loopback(cx, &mut client, &mut server).expect("loopback establishes");
        (client, server)
    }

    #[test]
    fn loopback_reaches_established_on_both_sides() {
        let cx = test_cx();
        let (client, server) = established_pair(&cx);
        assert_eq!(client.state(), QuicConnectionState::Established);
        assert_eq!(server.state(), QuicConnectionState::Established);
        assert!(client.can_send_app_data());
        assert!(server.can_send_app_data());
        assert_eq!(client.role(), StreamRole::Client);
        assert_eq!(server.role(), StreamRole::Server);
    }

    #[test]
    fn client_handshake_fails_closed_without_verified_identity() {
        let cx = test_cx();
        let (mut client, mut server) = pair();
        // No record_verified_server_identity() -> client confirm must fail closed.
        let err = establish_loopback(&cx, &mut client, &mut server)
            .expect_err("client must fail closed without verified identity");
        assert!(
            matches!(err, NativeQuicConnectionError::Tls(_)),
            "expected a TLS fail-closed error, got {err:?}"
        );
        assert_ne!(client.state(), QuicConnectionState::Established);
    }

    #[test]
    fn datagram_roundtrip_exact_and_fifo() {
        let cx = test_cx();
        let (mut client, mut server) = established_pair(&cx);

        let payloads: [&[u8]; 3] = [b"first-symbol", b"second", b"third-datagram-payload"];
        for p in payloads {
            client
                .send_datagram(&cx, Bytes::copy_from_slice(p))
                .expect("queue datagram");
        }
        let moved = pump_until_idle(
            &cx,
            &mut client,
            &mut server,
            DEFAULT_MAX_PACKET_BYTES,
            1000,
        )
        .expect("pump");
        assert!(
            moved >= payloads.len(),
            "expected >= {} frames",
            payloads.len()
        );

        for expected in payloads {
            let got = server.recv_datagram().expect("a datagram arrived");
            assert_eq!(got.as_ref(), expected, "exact payload, FIFO order");
        }
        assert!(server.recv_datagram().is_none(), "no extra datagrams");
        assert_eq!(client.datagrams_sent(), payloads.len() as u64);
        assert_eq!(server.datagrams_received(), payloads.len() as u64);
    }

    #[test]
    fn datagram_send_before_established_is_rejected() {
        let cx = test_cx();
        let (mut client, _server) = pair();
        let err = client
            .send_datagram(&cx, Bytes::from_static(b"too early"))
            .expect_err("must reject before established");
        assert!(matches!(err, NativeQuicConnectionError::InvalidState(_)));
    }

    #[test]
    fn oversize_datagram_is_rejected_fail_closed() {
        let cx = test_cx();
        let (mut client, _server) = established_pair(&cx);
        let huge = Bytes::from(vec![0xABu8; 4096]);
        let err = client
            .send_datagram(&cx, huge)
            .expect_err("oversize datagram must be rejected");
        assert!(matches!(
            err,
            NativeQuicConnectionError::DatagramTooLarge { .. }
        ));
    }

    #[test]
    fn control_stream_roundtrip_multi_packet_reassembly() {
        let cx = test_cx();
        let (mut client, mut server) = established_pair(&cx);

        let stream = client
            .open_control_stream(&cx)
            .expect("open control stream");
        // A payload large enough to span several small packets.
        let body: Vec<u8> = (0..2048u32).map(|i| (i % 251) as u8).collect();
        client
            .write_control(&cx, stream, Bytes::copy_from_slice(&body), true)
            .expect("write control bytes + FIN");

        // Tiny budget forces multi-packet fragmentation across the pump.
        let moved = pump_until_idle(&cx, &mut client, &mut server, 256, 2000).expect("pump");
        assert!(
            moved > 1,
            "large payload should span multiple packets, moved {moved}"
        );

        let mut received = Vec::new();
        loop {
            let chunk = server
                .read_control(&cx, stream, 4096)
                .expect("read control");
            if chunk.is_empty() {
                break;
            }
            received.extend_from_slice(&chunk);
        }
        assert_eq!(received, body, "stream bytes reassembled in order");
        assert!(
            server.is_control_eof(stream).expect("eof query"),
            "FIN should be observed after consuming all bytes"
        );
    }

    #[test]
    fn poll_recv_datagram_pending_then_ready() {
        use std::sync::Arc;
        use std::sync::atomic::{AtomicUsize, Ordering};
        use std::task::Wake;

        struct CountingWaker(AtomicUsize);
        impl Wake for CountingWaker {
            fn wake(self: Arc<Self>) {
                self.0.fetch_add(1, Ordering::SeqCst);
            }
            fn wake_by_ref(self: &Arc<Self>) {
                self.0.fetch_add(1, Ordering::SeqCst);
            }
        }

        let cx = test_cx();
        let (mut client, mut server) = established_pair(&cx);

        let counter = Arc::new(CountingWaker(AtomicUsize::new(0)));
        let waker = counter.clone().into();
        let mut task_cx = TaskContext::from_waker(&waker);

        // Empty queue: registers the waker, returns Pending (no busy-poll).
        assert!(matches!(
            server.poll_recv_datagram(&cx, &mut task_cx),
            Poll::Pending
        ));

        // Deliver a datagram; the registered waker must fire.
        client
            .send_datagram(&cx, Bytes::from_static(b"wakeup"))
            .expect("queue datagram");
        pump_until_idle(
            &cx,
            &mut client,
            &mut server,
            DEFAULT_MAX_PACKET_BYTES,
            3000,
        )
        .expect("pump");
        assert!(
            counter.0.load(Ordering::SeqCst) >= 1,
            "arrival must wake the registered task"
        );

        // Now the poll resolves with the exact payload.
        match server.poll_recv_datagram(&cx, &mut task_cx) {
            Poll::Ready(Ok(got)) => assert_eq!(got.as_ref(), b"wakeup"),
            other => panic!("expected Ready(Ok), got {other:?}"),
        }
    }

    #[test]
    fn path_stats_are_exposed() {
        let cx = test_cx();
        let (client, _server) = established_pair(&cx);
        let stats = client.path_stats();
        // Congestion window has a nonzero initial value; in-flight starts at 0.
        assert!(stats.congestion_window_bytes > 0);
        assert_eq!(stats.bytes_in_flight, 0);
        assert_eq!(stats.pto_count, 0);
        assert_eq!(stats.packets_acked, 0);
        assert_eq!(stats.packets_lost, 0);
        assert_eq!(stats.loss_rate, 0.0);
    }

    #[test]
    fn graceful_close_transitions_out_of_established() {
        let cx = test_cx();
        let (mut client, _server) = established_pair(&cx);
        client.begin_close(&cx, 5000, 0).expect("begin close");
        assert_ne!(client.state(), QuicConnectionState::Established);
    }
}