This commit is contained in:
Alexey
2026-09-25 19:46:31 +03:00
parent 2d63fcf376
commit 326c0ecdb9
128 changed files with 914 additions and 1248 deletions
+2 -5
View File
@@ -33,11 +33,8 @@ async fn stale_generation_cannot_publish_bootstrap_after_disabled_cutover() {
let disabled = test_runtime_generation(2, disabled_config);
runtime.activate_generation(Arc::clone(&disabled));
let result = runtime.issue_bootstrap_for_generation(
&initial,
profile,
"192.0.2.10".parse().unwrap(),
);
let result =
runtime.issue_bootstrap_for_generation(&initial, profile, "192.0.2.10".parse().unwrap());
assert!(matches!(result, Err(ManagerError::Closed)));
let status = serde_json::to_value(runtime.try_status()).unwrap();
+5 -3
View File
@@ -113,9 +113,11 @@ async fn user_revocation_interrupts_live_session_before_periodic_cleanup() {
.unwrap(),
Err(ManagerError::Closed)
));
assert!(runtime
.get_session(session_hash, "proxy.example.com")
.is_err());
assert!(
runtime
.get_session(session_hash, "proxy.example.com")
.is_err()
);
stop_runtime(runtime, generation).await;
}
+1 -1
View File
@@ -5,8 +5,8 @@ use bytes::Bytes;
use tokio_tungstenite::tungstenite::protocol::Message;
use tokio_util::sync::CancellationToken;
use super::{CarrierSocket, DataPlaneEvent, DriverEvent, FairDataSelector};
use super::io::{flush, process_lane, read_message, record_message, reserve_data, send};
use super::{CarrierSocket, DataPlaneEvent, DriverEvent, FairDataSelector};
use crate::web::manager::{WebProcessRuntime, WebSocketConnection};
use crate::web::session::{SessionCloseReason, WebSession, WebSocketLaneReservation};
use crate::web::trace::{TraceDirection, TraceWebSocketContext};
+1 -1
View File
@@ -56,8 +56,8 @@ mod observability;
pub(crate) use observability::{WebCapacityResourceStatus, WebCapacitySnapshot};
// Asynchronous bounded close operations isolate mutation lifecycle from HTTP requests.
mod control;
pub(crate) use budget::WebSocketBudgetLease;
use budget::WebDataBudget;
pub(crate) use budget::WebSocketBudgetLease;
pub(crate) use control::{CloseOperationSelector, ControlError};
pub(crate) use negotiation::{
CarrierCapabilities, CarrierClientClass, CarrierFailure, CarrierLearningContext, CarrierRequest,
+5 -12
View File
@@ -18,12 +18,8 @@ impl WebProcessRuntime {
client_ip: IpAddr,
public_addr: SocketAddr,
) -> Result<u16, super::ManagerError> {
let (result, notify) = self.try_acquire_stream_quiet(
profile_key,
max_streams,
client_ip,
public_addr,
);
let (result, notify) =
self.try_acquire_stream_quiet(profile_key, max_streams, client_ip, public_addr);
if let Some(notify) = notify {
notify.notify_waiters();
}
@@ -94,12 +90,9 @@ impl WebProcessRuntime {
public_addr: SocketAddr,
peer_port: u16,
) {
if let Some(notify) = self.release_stream_quiet(
profile_key,
client_ip,
public_addr,
peer_port,
) {
if let Some(notify) =
self.release_stream_quiet(profile_key, client_ip, public_addr, peer_port)
{
notify.notify_waiters();
}
}
+4 -1
View File
@@ -66,7 +66,10 @@ fn quiet_queue_release_updates_accounting_before_notification_dispatch() {
let waker = Waker::from(Arc::clone(&counter));
let mut context = Context::from_waker(&waker);
let mut notified = Box::pin(budget.notify.notified());
assert!(matches!(notified.as_mut().poll(&mut context), Poll::Pending));
assert!(matches!(
notified.as_mut().poll(&mut context),
Poll::Pending
));
let notify = budget.release_queue_quiet([1; 32], 64, 1, false);
+1 -3
View File
@@ -117,9 +117,7 @@ impl WebProcessRuntime {
return Err(ManagerError::Limit);
}
let global_capacity_full = state.bootstraps.len() >= self.limits.max_bootstraps_global;
if global_capacity_full
&& !state.bootstraps.values().any(|bootstrap| !bootstrap.used)
{
if global_capacity_full && !state.bootstraps.values().any(|bootstrap| !bootstrap.used) {
self.record_limit_hit();
self.telemetry
.record_rejection(WebRejectionReason::BootstrapCapacity);
@@ -207,7 +207,10 @@ mod tests {
let waker = Waker::from(Arc::clone(&counter));
let mut context = Context::from_waker(&waker);
let mut notified = Box::pin(admission.registrations_drained.notified());
assert!(matches!(notified.as_mut().poll(&mut context), Poll::Pending));
assert!(matches!(
notified.as_mut().poll(&mut context),
Poll::Pending
));
let notify = registration.release_deferred().unwrap();
+4 -1
View File
@@ -92,7 +92,10 @@ async fn quiet_stream_release_returns_exact_post_accounting_drain_notification()
let waker = Waker::from(Arc::clone(&counter));
let mut context = Context::from_waker(&waker);
let mut notified = Box::pin(runtime.operator_lifecycle.work_changed.notified());
assert!(matches!(notified.as_mut().poll(&mut context), Poll::Pending));
assert!(matches!(
notified.as_mut().poll(&mut context),
Poll::Pending
));
let notify = runtime
.release_stream_quiet(profile_key, client_ip, public_addr, peer_port)
@@ -40,18 +40,15 @@ impl WebProcessRuntime {
.sessions
.get(&replacement.old_session.token_hash())
.is_some_and(|session| Arc::ptr_eq(session, &replacement.old_session));
if !valid
|| state.closed
|| !state.issuance_enabled
{
if !valid || state.closed || !state.issuance_enabled {
drop(state);
self.cancel_replacement(bootstrap_hash, &replacement.old_session);
return Err(ManagerError::Closed);
}
let Some(mut user_publication) = generation.proxy_shared.claim_authenticated_user(
&replacement.profile.user,
replacement.profile.credential_id,
) else {
let Some(mut user_publication) = generation
.proxy_shared
.claim_authenticated_user(&replacement.profile.user, replacement.profile.credential_id)
else {
drop(state);
self.cancel_replacement(bootstrap_hash, &replacement.old_session);
return Err(ManagerError::Closed);
+1 -2
View File
@@ -11,11 +11,11 @@ use tokio::sync::Notify;
use tokio_util::sync::CancellationToken;
use crate::config::{WebCarrier, WebLimitsConfig, WebRuntimeProfile, WebTimeoutsConfig};
use crate::proxy::user_admission::UserSessionRegistration;
use crate::web::frame::FrameType;
use crate::web::manager::{
CarrierClientClass, CarrierLearningContext, ProfileKey, TokenHash, WebProcessRuntime,
};
use crate::proxy::user_admission::UserSessionRegistration;
// Backend tasks own generation admission and authenticated MTProxy relay lifetimes.
mod backend;
@@ -416,5 +416,4 @@ impl WebSession {
pub(crate) fn timeouts(&self) -> &WebTimeoutsConfig {
&self.timeouts
}
}
+16 -38
View File
@@ -116,21 +116,18 @@ impl WebSession {
});
}
if !state.pending_frames.is_empty() {
let batch = match self.take_down_batch_locked(
&mut state,
&mut effects,
cursor,
) {
Ok(batch) => batch,
Err(ManagerError::Backpressure) => {
return Err(ManagerError::Backpressure);
}
Err(error) => {
drop(state);
self.close(SessionCloseReason::Protocol);
return Err(error);
}
};
let batch =
match self.take_down_batch_locked(&mut state, &mut effects, cursor) {
Ok(batch) => batch,
Err(ManagerError::Backpressure) => {
return Err(ManagerError::Backpressure);
}
Err(error) => {
drop(state);
self.close(SessionCloseReason::Protocol);
return Err(error);
}
};
let result = PollResult {
body: batch.body.clone(),
next_cursor: batch.next_cursor,
@@ -300,12 +297,7 @@ impl WebSession {
state.pending_control_items = state.pending_control_items.saturating_sub(items);
}
if let Some(manager) = self.manager.upgrade() {
effects.notify(manager.release_pending_quiet(
self.profile_key,
bytes,
items,
control,
));
effects.notify(manager.release_pending_quiet(self.profile_key, bytes, items, control));
}
}
@@ -417,14 +409,7 @@ impl WebSession {
last.encoded[4..8].copy_from_slice(&payload_len.to_be_bytes());
return true;
}
self.queue_frame_locked(
state,
effects,
FrameType::Data,
stream_id,
payload,
false,
)
self.queue_frame_locked(state, effects, FrameType::Data, stream_id, payload, false)
}
fn queue_frame_locked(
@@ -437,14 +422,8 @@ impl WebSession {
control: bool,
) -> bool {
if self.carrier().uses_lanes() {
return self.queue_lane_frame_locked(
state,
effects,
frame_type,
stream_id,
payload,
control,
);
return self
.queue_lane_frame_locked(state, effects, frame_type, stream_id, payload, control);
}
let cost = frame::HEADER_BYTES + payload.len() + QUEUE_ITEM_COST;
let class = if control {
@@ -476,7 +455,6 @@ impl WebSession {
effects.notify(Arc::clone(&self.down_notify));
true
}
}
#[cfg(test)]
+8 -2
View File
@@ -89,7 +89,10 @@ async fn queued_frame_notifies_only_after_releasing_session_lock() {
}));
let mut context = Context::from_waker(&waker);
let mut notified = Box::pin(session.down_notify.notified());
assert!(matches!(notified.as_mut().poll(&mut context), Poll::Pending));
assert!(matches!(
notified.as_mut().poll(&mut context),
Poll::Pending
));
queue_close(&session);
@@ -109,7 +112,10 @@ async fn budget_release_notifies_only_after_session_accounting_and_unlock() {
}));
let mut context = Context::from_waker(&waker);
let mut notified = Box::pin(manager.budget_notify().notified_owned());
assert!(matches!(notified.as_mut().poll(&mut context), Poll::Pending));
assert!(matches!(
notified.as_mut().poll(&mut context),
Poll::Pending
));
session.with_state_effects(|state, effects| {
let bytes = state.pending_control_bytes;
+4 -8
View File
@@ -78,8 +78,7 @@ impl DeferredSessionEffects {
/// Defers one RawWaker drop without delivering a readiness signal.
pub(super) fn drop_waker(&mut self, waker: Waker) {
self.callbacks
.push(DeferredSessionEffect::DropWaker(waker));
self.callbacks.push(DeferredSessionEffect::DropWaker(waker));
}
/// Defers one exact notification without coalescing sibling effects.
@@ -89,8 +88,7 @@ impl DeferredSessionEffects {
/// Retains a detached response batch until its lease can drop safely.
pub(super) fn retain_batch(&mut self, batch: DownBatch) {
self.retained
.push(RetainedSessionResource::Batch(batch));
self.retained.push(RetainedSessionResource::Batch(batch));
}
/// Retains transient semaphore capacity until state publication completes.
@@ -173,10 +171,8 @@ mod tests {
impl Wake for PermitOrderWake {
fn wake(self: Arc<Self>) {
self.observed_release.store(
self.semaphore.available_permits(),
Ordering::Release,
);
self.observed_release
.store(self.semaphore.available_permits(), Ordering::Release);
}
}
+1 -7
View File
@@ -181,13 +181,7 @@ impl WebSession {
&mut unused_items,
&mut progress,
);
self.release_locked(
&mut state,
&mut effects,
unused_bytes,
unused_items,
false,
);
self.release_locked(&mut state, &mut effects, unused_bytes, unused_items, false);
if let Some(lane) = state.carrier_lanes.get_mut(&lane_id) {
lane.up_active = false;
if applied {
+6 -11
View File
@@ -190,12 +190,11 @@ impl WebSession {
{
return None;
}
let released =
self.release_on_close_locked(
&mut state,
SessionCloseReason::CarrierSuperseded,
effects,
);
let released = self.release_on_close_locked(
&mut state,
SessionCloseReason::CarrierSuperseded,
effects,
);
Some(CarrierSupersedeCompletion {
session: self,
released,
@@ -264,11 +263,7 @@ impl WebSession {
{
return None;
}
Some(self.release_on_close_locked(
&mut state,
SessionCloseReason::PeerIdle,
effects,
))
Some(self.release_on_close_locked(&mut state, SessionCloseReason::PeerIdle, effects))
}
fn release_on_close_locked(
+2 -3
View File
@@ -1,8 +1,8 @@
use std::collections::VecDeque;
use std::mem::ManuallyDrop;
use std::net::SocketAddr;
use std::sync::{Arc, Barrier, Weak};
use std::sync::atomic::{AtomicBool, AtomicU8, AtomicUsize, Ordering};
use std::sync::{Arc, Barrier, Weak};
use std::task::{Context, Poll, RawWaker, RawWakerVTable, Waker};
use bytes::Bytes;
@@ -63,8 +63,7 @@ impl CallbackProbe {
self.clone_while_locked.store(true, Ordering::Release);
return;
};
if self.clone_action.swap(CLONE_ACTION_NONE, Ordering::AcqRel)
== CLONE_ACTION_INSERT_DATA
if self.clone_action.swap(CLONE_ACTION_NONE, Ordering::AcqRel) == CLONE_ACTION_INSERT_DATA
&& let Some(stream) = state
.streams
.get_mut(&self.stream.id)
+1 -7
View File
@@ -149,13 +149,7 @@ impl WebSession {
&mut unused_items,
&mut progress,
);
self.release_locked(
&mut state,
&mut effects,
unused_bytes,
unused_items,
false,
);
self.release_locked(&mut state, &mut effects, unused_bytes, unused_items, false);
if !applied {
Err(ManagerError::Closed)
} else {
+4 -18
View File
@@ -14,8 +14,8 @@ use crate::web::manager::ManagerError;
// Reservation ownership keeps pre-OPEN quota and exact lane identity transactional.
mod reservation;
pub(crate) use reservation::{WebSocketLaneReservation, WebSocketProbeReservation};
use reservation::WebSocketLaneReservationPhase;
pub(crate) use reservation::{WebSocketLaneReservation, WebSocketProbeReservation};
impl WebSession {
/// Reserves the only automatic WebSocket probe before any HTTP 101 response.
@@ -268,13 +268,7 @@ impl WebSession {
&mut unused_items,
&mut progress,
);
self.release_locked(
&mut state,
&mut effects,
unused_bytes,
unused_items,
false,
);
self.release_locked(&mut state, &mut effects, unused_bytes, unused_items, false);
if let Some(lane) = state.carrier_lanes.get_mut(&lane_id) {
lane.up_active = false;
if applied {
@@ -385,21 +379,13 @@ impl WebSession {
state.active_peer_ports.remove(&claim.peer_port)
};
if lane_matches {
self.remember_closed_locked(
&mut state,
&mut effects,
claim.lane.lane_id,
);
self.remember_closed_locked(&mut state, &mut effects, claim.lane.lane_id);
if state
.carrier_lanes
.get(&claim.lane.lane_id)
.is_some_and(|lane| lane.instance == claim.lane.instance)
{
self.release_lane_locked(
&mut state,
&mut effects,
claim.lane.lane_id,
);
self.release_lane_locked(&mut state, &mut effects, claim.lane.lane_id);
}
}
release_port
+1 -4
View File
@@ -156,10 +156,7 @@ impl WebSocketLaneReservation {
Ok(())
}
pub(super) fn mark_stream_owned(
&mut self,
stream: StreamIdentity,
) -> Result<(), ManagerError> {
pub(super) fn mark_stream_owned(&mut self, stream: StreamIdentity) -> Result<(), ManagerError> {
if self.phase != WebSocketLaneReservationPhase::Transferred || self.stream != Some(stream) {
return Err(ManagerError::Protocol);
}
+1 -3
View File
@@ -324,9 +324,7 @@ impl HttpTraceExchange {
.clamp(1, limits.max_frames_per_body);
let reservation = estimated_frames.saturating_mul(std::mem::size_of::<TraceFrame>());
let mut state = self.state.lock();
if state.phase != ExchangePhase::Open
|| !self.reserve_locked(&mut state, reservation)
{
if state.phase != ExchangePhase::Open || !self.reserve_locked(&mut state, reservation) {
return;
}
let frames = match frame::parse_all(body, limits) {