Loading src/quic.cpp +15 −2 Original line number Diff line number Diff line Loading @@ -120,6 +120,7 @@ void quicPerfReportOnce() { static uint64_t last_sb_flushes = 0, last_sb_packets = 0; static uint64_t last_rb_calls = 0, last_rb_packets = 0; static uint64_t last_pumps_done = 0, last_pumps_skipped = 0; static uint64_t last_cc_spin_iter = 0, last_cc_wait_iter = 0; auto& perf = quicPerf(); auto ld = [](std::atomic<uint64_t>& a) { return a.load(std::memory_order_relaxed); }; Loading @@ -135,6 +136,7 @@ void quicPerfReportOnce() { uint64_t sb_flushes = ld(perf.send_batch_flushes), sb_packets = ld(perf.send_batch_packets); uint64_t rb_calls = ld(perf.recv_batch_calls), rb_packets = ld(perf.recv_batch_packets); uint64_t pumps_done = ld(perf.pumps_done), pumps_skipped = ld(perf.pumps_skipped); uint64_t cc_spin_iter = ld(perf.cc_spin_iterations), cc_wait_iter = ld(perf.cc_wait_iterations); uint64_t d_bytes_sent = bytes_sent - last_bytes_sent; uint64_t d_bytes_recv = bytes_recv - last_bytes_recv; Loading @@ -158,6 +160,8 @@ void quicPerfReportOnce() { uint64_t d_rb_packets = rb_packets - last_rb_packets; uint64_t d_pumps_done = pumps_done - last_pumps_done; uint64_t d_pumps_skipped = pumps_skipped - last_pumps_skipped; uint64_t d_cc_spin_iter = cc_spin_iter - last_cc_spin_iter; uint64_t d_cc_wait_iter = cc_wait_iter - last_cc_wait_iter; double tx_mb = double(d_bytes_sent) / (1024.0 * 1024.0); double rx_mb = double(d_bytes_recv) / (1024.0 * 1024.0); Loading @@ -167,12 +171,14 @@ void quicPerfReportOnce() { // Approximation, not an exact RFC 9002 loss rate: lost-packets-detected // divided by packets sent this interval. Good enough to see the trend. double loss_pct = d_packets_sent ? (100.0 * double(d_lost) / double(d_packets_sent)) : 0.0; double avg_cc_spin = d_cc_stalls ? double(d_cc_spin_iter) / double(d_cc_stalls) : 0.0; double avg_cc_wait = d_cc_stalls ? double(d_cc_wait_iter) / double(d_cc_stalls) : 0.0; std::fprintf(stderr, "[QUIC_PERF] tx=%.2fMB/s rx=%.2fMB/s pkt_tx=%lu/s pkt_rx=%lu/s | " "gso=%lu/s segs/gso=%.1f sendmmsg=%lu/s single_send=%lu/s recvmmsg=%lu/s gro=%lu/s | " "cwnd=%lu inflight=%lu ssthresh=%lu rtt=%.2fms minrtt=%.2fms | " "loss=%.2f%% retx=%lu/s eagain(tx/rx)=%lu/%lu fcstall=%lu/s ccstall=%lu/s | " "loss=%.2f%% retx=%lu/s eagain(tx/rx)=%lu/%lu fcstall=%lu/s ccstall=%lu/s ccspin/wait=%.1f/%.1f | " "sendbatch(avg/max)=%.1f/%lu recvbatch(avg/max)=%.1f/%lu pumps(done/skip)=%lu/%lu\n", tx_mb, rx_mb, (unsigned long)d_packets_sent, (unsigned long)d_packets_recv, Loading @@ -182,7 +188,7 @@ void quicPerfReportOnce() { ld(perf.last_rtt_us) / 1000.0, ld(perf.last_min_rtt_us) / 1000.0, loss_pct, (unsigned long)d_retransmits, (unsigned long)d_send_eagain, (unsigned long)d_recv_eagain, (unsigned long)d_fc_stalls, (unsigned long)d_cc_stalls, (unsigned long)d_fc_stalls, (unsigned long)d_cc_stalls, avg_cc_spin, avg_cc_wait, avg_send_batch, (unsigned long)ld(perf.send_batch_max), avg_recv_batch, (unsigned long)ld(perf.recv_batch_max), (unsigned long)d_pumps_done, (unsigned long)d_pumps_skipped); Loading @@ -199,6 +205,7 @@ void quicPerfReportOnce() { last_sb_flushes = sb_flushes; last_sb_packets = sb_packets; last_rb_calls = rb_calls; last_rb_packets = rb_packets; last_pumps_done = pumps_done; last_pumps_skipped = pumps_skipped; last_cc_spin_iter = cc_spin_iter; last_cc_wait_iter = cc_wait_iter; } void quicPerfReporterLoop() { Loading Loading @@ -7807,6 +7814,7 @@ size_t quic::sendStreamData(uint64_t stream_id, const uint8_t* data, size_t len, // Drain incoming to receive ACKs and free cwnd // First: spin-pump a few times without syscall (ACK may already be queued) bool unblocked = false; int spin_count = 0; for (int spin = 0; spin < 8; ++spin) { // quic_mtx() is already held throughout this function (see // `lock` above) — pumpIncomingLocked() skips the redundant Loading @@ -7814,8 +7822,10 @@ size_t quic::sendStreamData(uint64_t stream_id, const uint8_t* data, size_t len, // up to 8 times per congestion-window stall (F31). pumpIncomingLocked(); _flushes_since_pump = 0; ++spin_count; if (cwndAllowsSend(est_pkt_size)) { unblocked = true; break; } } quicPerf().cc_spin_iterations.fetch_add(static_cast<uint64_t>(spin_count), std::memory_order_relaxed); // If still blocked, poll with timeout. Reuse one socketwait // (and its epoll fd) across every retry instead of constructing // a fresh one per ~1ms attempt — see _cc_socketwait's comment. Loading @@ -7832,6 +7842,7 @@ size_t quic::sendStreamData(uint64_t stream_id, const uint8_t* data, size_t len, // estimate can't stall the sender indefinitely. int cc_wait_budget_ms = static_cast<int>( std::clamp(_rto * 1000.0, 100.0, 2000.0)); int wait_count = 0; for (int cc_wait = 0; cc_wait < cc_wait_budget_ms; ++cc_wait) { if (_conn_state.load() == ConnectionState::Closed || _conn_state.load() == ConnectionState::Draining) break; Loading @@ -7842,8 +7853,10 @@ size_t quic::sendStreamData(uint64_t stream_id, const uint8_t* data, size_t len, // comment on pumpIncomingLocked() (F31). pumpIncomingLocked(); _flushes_since_pump = 0; ++wait_count; if (cwndAllowsSend(est_pkt_size)) { unblocked = true; break; } } quicPerf().cc_wait_iterations.fetch_add(static_cast<uint64_t>(wait_count), std::memory_order_relaxed); } if (!unblocked) break; } Loading src/socket.h +6 −0 Original line number Diff line number Diff line Loading @@ -612,6 +612,12 @@ namespace netplus { std::atomic<uint64_t> retransmits{0}; std::atomic<uint64_t> flow_control_stalls{0}; // DATA_BLOCKED/STREAM_DATA_BLOCKED sent std::atomic<uint64_t> congestion_stalls{0}; // cwndAllowsSend() block entries // Diagnostic-only: how many of the 8 non-blocking spin iterations // and how many of the ≤cc_wait_budget_ms timed-wait iterations // sendStreamData()'s congestion-block path actually consumed per // stall — see its own comment. avg = sum/congestion_stalls. std::atomic<uint64_t> cc_spin_iterations{0}; std::atomic<uint64_t> cc_wait_iterations{0}; std::atomic<uint64_t> pumps_done{0}; // pumpIncomingLocked() invocations from the send loop std::atomic<uint64_t> pumps_skipped{0}; // adaptive-pump skips (had cwnd headroom) Loading Loading
src/quic.cpp +15 −2 Original line number Diff line number Diff line Loading @@ -120,6 +120,7 @@ void quicPerfReportOnce() { static uint64_t last_sb_flushes = 0, last_sb_packets = 0; static uint64_t last_rb_calls = 0, last_rb_packets = 0; static uint64_t last_pumps_done = 0, last_pumps_skipped = 0; static uint64_t last_cc_spin_iter = 0, last_cc_wait_iter = 0; auto& perf = quicPerf(); auto ld = [](std::atomic<uint64_t>& a) { return a.load(std::memory_order_relaxed); }; Loading @@ -135,6 +136,7 @@ void quicPerfReportOnce() { uint64_t sb_flushes = ld(perf.send_batch_flushes), sb_packets = ld(perf.send_batch_packets); uint64_t rb_calls = ld(perf.recv_batch_calls), rb_packets = ld(perf.recv_batch_packets); uint64_t pumps_done = ld(perf.pumps_done), pumps_skipped = ld(perf.pumps_skipped); uint64_t cc_spin_iter = ld(perf.cc_spin_iterations), cc_wait_iter = ld(perf.cc_wait_iterations); uint64_t d_bytes_sent = bytes_sent - last_bytes_sent; uint64_t d_bytes_recv = bytes_recv - last_bytes_recv; Loading @@ -158,6 +160,8 @@ void quicPerfReportOnce() { uint64_t d_rb_packets = rb_packets - last_rb_packets; uint64_t d_pumps_done = pumps_done - last_pumps_done; uint64_t d_pumps_skipped = pumps_skipped - last_pumps_skipped; uint64_t d_cc_spin_iter = cc_spin_iter - last_cc_spin_iter; uint64_t d_cc_wait_iter = cc_wait_iter - last_cc_wait_iter; double tx_mb = double(d_bytes_sent) / (1024.0 * 1024.0); double rx_mb = double(d_bytes_recv) / (1024.0 * 1024.0); Loading @@ -167,12 +171,14 @@ void quicPerfReportOnce() { // Approximation, not an exact RFC 9002 loss rate: lost-packets-detected // divided by packets sent this interval. Good enough to see the trend. double loss_pct = d_packets_sent ? (100.0 * double(d_lost) / double(d_packets_sent)) : 0.0; double avg_cc_spin = d_cc_stalls ? double(d_cc_spin_iter) / double(d_cc_stalls) : 0.0; double avg_cc_wait = d_cc_stalls ? double(d_cc_wait_iter) / double(d_cc_stalls) : 0.0; std::fprintf(stderr, "[QUIC_PERF] tx=%.2fMB/s rx=%.2fMB/s pkt_tx=%lu/s pkt_rx=%lu/s | " "gso=%lu/s segs/gso=%.1f sendmmsg=%lu/s single_send=%lu/s recvmmsg=%lu/s gro=%lu/s | " "cwnd=%lu inflight=%lu ssthresh=%lu rtt=%.2fms minrtt=%.2fms | " "loss=%.2f%% retx=%lu/s eagain(tx/rx)=%lu/%lu fcstall=%lu/s ccstall=%lu/s | " "loss=%.2f%% retx=%lu/s eagain(tx/rx)=%lu/%lu fcstall=%lu/s ccstall=%lu/s ccspin/wait=%.1f/%.1f | " "sendbatch(avg/max)=%.1f/%lu recvbatch(avg/max)=%.1f/%lu pumps(done/skip)=%lu/%lu\n", tx_mb, rx_mb, (unsigned long)d_packets_sent, (unsigned long)d_packets_recv, Loading @@ -182,7 +188,7 @@ void quicPerfReportOnce() { ld(perf.last_rtt_us) / 1000.0, ld(perf.last_min_rtt_us) / 1000.0, loss_pct, (unsigned long)d_retransmits, (unsigned long)d_send_eagain, (unsigned long)d_recv_eagain, (unsigned long)d_fc_stalls, (unsigned long)d_cc_stalls, (unsigned long)d_fc_stalls, (unsigned long)d_cc_stalls, avg_cc_spin, avg_cc_wait, avg_send_batch, (unsigned long)ld(perf.send_batch_max), avg_recv_batch, (unsigned long)ld(perf.recv_batch_max), (unsigned long)d_pumps_done, (unsigned long)d_pumps_skipped); Loading @@ -199,6 +205,7 @@ void quicPerfReportOnce() { last_sb_flushes = sb_flushes; last_sb_packets = sb_packets; last_rb_calls = rb_calls; last_rb_packets = rb_packets; last_pumps_done = pumps_done; last_pumps_skipped = pumps_skipped; last_cc_spin_iter = cc_spin_iter; last_cc_wait_iter = cc_wait_iter; } void quicPerfReporterLoop() { Loading Loading @@ -7807,6 +7814,7 @@ size_t quic::sendStreamData(uint64_t stream_id, const uint8_t* data, size_t len, // Drain incoming to receive ACKs and free cwnd // First: spin-pump a few times without syscall (ACK may already be queued) bool unblocked = false; int spin_count = 0; for (int spin = 0; spin < 8; ++spin) { // quic_mtx() is already held throughout this function (see // `lock` above) — pumpIncomingLocked() skips the redundant Loading @@ -7814,8 +7822,10 @@ size_t quic::sendStreamData(uint64_t stream_id, const uint8_t* data, size_t len, // up to 8 times per congestion-window stall (F31). pumpIncomingLocked(); _flushes_since_pump = 0; ++spin_count; if (cwndAllowsSend(est_pkt_size)) { unblocked = true; break; } } quicPerf().cc_spin_iterations.fetch_add(static_cast<uint64_t>(spin_count), std::memory_order_relaxed); // If still blocked, poll with timeout. Reuse one socketwait // (and its epoll fd) across every retry instead of constructing // a fresh one per ~1ms attempt — see _cc_socketwait's comment. Loading @@ -7832,6 +7842,7 @@ size_t quic::sendStreamData(uint64_t stream_id, const uint8_t* data, size_t len, // estimate can't stall the sender indefinitely. int cc_wait_budget_ms = static_cast<int>( std::clamp(_rto * 1000.0, 100.0, 2000.0)); int wait_count = 0; for (int cc_wait = 0; cc_wait < cc_wait_budget_ms; ++cc_wait) { if (_conn_state.load() == ConnectionState::Closed || _conn_state.load() == ConnectionState::Draining) break; Loading @@ -7842,8 +7853,10 @@ size_t quic::sendStreamData(uint64_t stream_id, const uint8_t* data, size_t len, // comment on pumpIncomingLocked() (F31). pumpIncomingLocked(); _flushes_since_pump = 0; ++wait_count; if (cwndAllowsSend(est_pkt_size)) { unblocked = true; break; } } quicPerf().cc_wait_iterations.fetch_add(static_cast<uint64_t>(wait_count), std::memory_order_relaxed); } if (!unblocked) break; } Loading
src/socket.h +6 −0 Original line number Diff line number Diff line Loading @@ -612,6 +612,12 @@ namespace netplus { std::atomic<uint64_t> retransmits{0}; std::atomic<uint64_t> flow_control_stalls{0}; // DATA_BLOCKED/STREAM_DATA_BLOCKED sent std::atomic<uint64_t> congestion_stalls{0}; // cwndAllowsSend() block entries // Diagnostic-only: how many of the 8 non-blocking spin iterations // and how many of the ≤cc_wait_budget_ms timed-wait iterations // sendStreamData()'s congestion-block path actually consumed per // stall — see its own comment. avg = sum/congestion_stalls. std::atomic<uint64_t> cc_spin_iterations{0}; std::atomic<uint64_t> cc_wait_iterations{0}; std::atomic<uint64_t> pumps_done{0}; // pumpIncomingLocked() invocations from the send loop std::atomic<uint64_t> pumps_skipped{0}; // adaptive-pump skips (had cwnd headroom) Loading