Process-wide concurrency + Cancellation ownership fixes

This commit is contained in:
Alexey
2026-09-23 22:30:51 +03:00
parent baa9bfbb01
commit f1107c21d9
47 changed files with 2177 additions and 762 deletions
+70 -11
View File
@@ -1,3 +1,4 @@
use std::future::Future;
use std::sync::Arc;
use std::time::{Duration, Instant};
@@ -112,6 +113,40 @@ pub(super) async fn run_upgraded(
type CarrierSocket = WebSocketStream<ConnectionIo>;
enum DataPlaneEvent<I, D> {
Incoming(I),
Down(D),
}
#[derive(Default)]
struct FairDataSelector {
prefer_down: bool,
}
impl FairDataSelector {
async fn select<I, D>(&mut self, incoming: I, down: D) -> DataPlaneEvent<I::Output, D::Output>
where
I: Future,
D: Future,
{
let event = if self.prefer_down {
tokio::select! {
biased;
down = down => DataPlaneEvent::Down(down),
incoming = incoming => DataPlaneEvent::Incoming(incoming),
}
} else {
tokio::select! {
biased;
incoming = incoming => DataPlaneEvent::Incoming(incoming),
down = down => DataPlaneEvent::Down(down),
}
};
self.prefer_down = matches!(event, DataPlaneEvent::Incoming(_));
event
}
}
async fn run_multiplex(
socket: &mut CarrierSocket,
runtime: &Arc<WebProcessRuntime>,
@@ -134,26 +169,31 @@ async fn run_multiplex(
let write_timeout = Duration::from_secs(session.timeouts().websocket_write_secs);
let maximum_message = session.limits().carrier_batch_bytes;
let mut active = false;
let mut data_selector = FairDataSelector::default();
loop {
let down = session.poll_down_websocket(cursor);
tokio::pin!(down);
let event = tokio::select! {
biased;
_ = cancellation.cancelled() => return Err(()),
_ = tokio::time::sleep_until(open_deadline.into()), if !active => return Err(()),
_ = tokio::time::sleep_until(next_ping.into()) => DriverEvent::Liveness,
incoming = read_message(
socket,
runtime,
session.profile_key(),
&cancellation,
&mut read_budget,
maximum_message,
backpressure_timeout,
data = data_selector.select(
read_message(
socket,
runtime,
session.profile_key(),
&cancellation,
&mut read_budget,
maximum_message,
backpressure_timeout,
),
down,
) => {
DriverEvent::Incoming(incoming?)
match data {
DataPlaneEvent::Incoming(incoming) => DriverEvent::Incoming(incoming?),
DataPlaneEvent::Down(down) => DriverEvent::Down(down.map_err(|_| ())?),
}
}
down = &mut down => DriverEvent::Down(down.map_err(|_| ())?),
};
match event {
DriverEvent::Incoming((message, _budget)) => match message {
@@ -363,3 +403,22 @@ enum DriverEvent {
Down(crate::web::session::PollResult),
Liveness,
}
#[cfg(test)]
mod fairness_tests {
use super::*;
#[tokio::test]
async fn continuously_ready_directions_alternate() {
let mut selector = FairDataSelector::default();
for expected_incoming in [true, false, true, false] {
let event = selector
.select(std::future::ready("incoming"), std::future::ready("down"))
.await;
assert_eq!(
matches!(event, DataPlaneEvent::Incoming(_)),
expected_incoming
);
}
}
}
+18 -19
View File
@@ -5,9 +5,9 @@ use bytes::Bytes;
use tokio_tungstenite::tungstenite::protocol::Message;
use tokio_util::sync::CancellationToken;
use super::CarrierSocket;
use super::{CarrierSocket, DataPlaneEvent, DriverEvent, FairDataSelector};
use super::io::{flush, process_lane, read_message, record_message, reserve_data, send};
use crate::web::manager::{WebProcessRuntime, WebSocketBudgetLease, WebSocketConnection};
use crate::web::manager::{WebProcessRuntime, WebSocketConnection};
use crate::web::session::{SessionCloseReason, WebSession, WebSocketLaneReservation};
use crate::web::trace::{TraceDirection, TraceWebSocketContext};
@@ -35,26 +35,31 @@ pub(super) async fn run_lane(
let write_timeout = Duration::from_secs(session.timeouts().websocket_write_secs);
let maximum_message = session.limits().carrier_batch_bytes;
let mut active = false;
let mut data_selector = FairDataSelector::default();
loop {
let down = session.poll_down_websocket_lane(reservation.lane_identity(), cursor);
tokio::pin!(down);
let event = tokio::select! {
biased;
_ = cancellation.cancelled() => return Err(()),
_ = tokio::time::sleep_until(open_deadline.into()), if !active => return Err(()),
_ = tokio::time::sleep_until(next_ping.into()) => DriverEvent::Liveness,
incoming = read_message(
socket,
runtime,
session.profile_key(),
&cancellation,
&mut read_budget,
maximum_message,
backpressure_timeout,
data = data_selector.select(
read_message(
socket,
runtime,
session.profile_key(),
&cancellation,
&mut read_budget,
maximum_message,
backpressure_timeout,
),
down,
) => {
DriverEvent::Incoming(incoming?)
match data {
DataPlaneEvent::Incoming(incoming) => DriverEvent::Incoming(incoming?),
DataPlaneEvent::Down(down) => DriverEvent::Down(down.map_err(|_| ())?),
}
}
down = &mut down => DriverEvent::Down(down.map_err(|_| ())?),
};
match event {
DriverEvent::Incoming((message, _budget)) => match message {
@@ -267,9 +272,3 @@ pub(super) async fn run_lane(
}
}
}
enum DriverEvent {
Incoming((Message, Option<WebSocketBudgetLease>)),
Down(crate::web::session::PollResult),
Liveness,
}
+19 -3
View File
@@ -1,6 +1,10 @@
use std::sync::Arc;
use std::sync::atomic::Ordering;
use std::task::Waker;
use std::time::{Duration, Instant};
use tokio::sync::Notify;
use super::WebSession;
/// Stable terminal cause assigned by the first session-close winner.
@@ -104,6 +108,8 @@ struct ReleasedQueues {
recovery_closed_before_commit: bool,
reason: SessionCloseReason,
peer_gap: Duration,
stream_wakers: Vec<Waker>,
lane_notifies: Vec<Arc<Notify>>,
}
/// Deferred queue release after manager publication linearizes a supersede.
@@ -276,12 +282,13 @@ impl WebSession {
if reason == SessionCloseReason::CarrierSuperseded {
state.negotiation_phase = SessionNegotiationPhase::Superseded;
}
let mut stream_wakers = Vec::with_capacity(state.streams.len().saturating_mul(2));
for stream in state.streams.values_mut() {
if let Some(waker) = stream.read_waker.take() {
waker.wake();
stream_wakers.push(waker);
}
if let Some(waker) = stream.write_waker.take() {
waker.wake();
stream_wakers.push(waker);
}
}
state.streams.clear();
@@ -296,8 +303,9 @@ impl WebSession {
let mut lane_data_items = 0usize;
let mut lane_control_bytes = 0usize;
let mut lane_control_items = 0usize;
let mut lane_notifies = Vec::with_capacity(state.carrier_lanes.len());
for lane in state.carrier_lanes.values_mut() {
lane.notify.notify_waiters();
lane_notifies.push(Arc::clone(&lane.notify));
if let Some(batch) = lane.unacked.take() {
batch.lease.detach();
lane_data_bytes = lane_data_bytes.saturating_add(batch.data_bytes);
@@ -326,10 +334,18 @@ impl WebSession {
recovery_closed_before_commit,
reason,
peer_gap,
stream_wakers,
lane_notifies,
}
}
fn finish_close(&self, released: ReleasedQueues) {
for waker in released.stream_wakers {
waker.wake();
}
for notify in released.lane_notifies {
notify.notify_waiters();
}
self.cancel.cancel();
if self.carrier().is_multiplexed() {
self.down_notify.notify_waiters();
+77
View File
@@ -1,10 +1,13 @@
use super::*;
use std::net::SocketAddr;
use std::sync::atomic::{AtomicBool, Ordering as AtomicOrdering};
use std::task::{Wake, Waker};
use crate::config::{
WebCarrier, WebLimitsConfig, WebRuntimeProfile, WebSecretMode, WebTimeoutsConfig,
};
use crate::web::manager::WebProcessRuntime;
use crate::web::session::SessionCloseOutcome;
fn session() -> Arc<WebSession> {
session_with_automatic(false)
@@ -53,6 +56,80 @@ fn session_with_automatic(automatic: bool) -> Arc<WebSession> {
)
}
struct SessionLockProbe {
session: std::sync::Weak<WebSession>,
lock_was_free: Arc<AtomicBool>,
}
impl Wake for SessionLockProbe {
fn wake(self: Arc<Self>) {
if let Some(session) = self.session.upgrade() {
self.lock_was_free
.store(session.state.try_lock().is_some(), AtomicOrdering::Release);
}
}
}
#[test]
fn close_wakes_stream_only_after_releasing_session_lock() {
let session = session();
let lock_was_free = Arc::new(AtomicBool::new(false));
let waker = Waker::from(Arc::new(SessionLockProbe {
session: Arc::downgrade(&session),
lock_was_free: Arc::clone(&lock_was_free),
}));
{
let mut state = session.state.lock();
state.streams.insert(
1,
StreamState {
instance: 1,
inbound: VecDeque::new(),
receive_window: frame::INITIAL_STREAM_WINDOW,
send_credit: u64::from(frame::INITIAL_STREAM_WINDOW),
read_waker: Some(waker),
write_waker: None,
},
);
}
assert_eq!(
session.close(SessionCloseReason::ApiClose),
SessionCloseOutcome::Closed
);
assert!(lock_was_free.load(AtomicOrdering::Acquire));
}
#[test]
fn supersede_completion_defers_stream_wake_until_finish() {
let session = session();
let lock_was_free = Arc::new(AtomicBool::new(false));
let waker = Waker::from(Arc::new(SessionLockProbe {
session: Arc::downgrade(&session),
lock_was_free: Arc::clone(&lock_was_free),
}));
{
let mut state = session.state.lock();
state.streams.insert(
1,
StreamState {
instance: 1,
inbound: VecDeque::new(),
receive_window: frame::INITIAL_STREAM_WINDOW,
send_credit: u64::from(frame::INITIAL_STREAM_WINDOW),
read_waker: Some(waker),
write_waker: None,
},
);
}
assert!(session.begin_carrier_supersede());
let completion = session.prepare_carrier_supersede().unwrap();
assert!(!lock_was_free.load(AtomicOrdering::Acquire));
completion.finish();
assert!(lock_was_free.load(AtomicOrdering::Acquire));
}
#[test]
fn uplink_retry_commits_only_one_exact_body() {
let session = session();