This commit is contained in:
Alexey
2026-07-13 12:20:24 +03:00
parent feb51cbf57
commit 73afeccae1
14 changed files with 105 additions and 151 deletions
@@ -145,8 +145,7 @@ me_writer_byte_budget_bytes = 268451840
"#, "#,
); );
let err = let err = ProxyConfig::load(&path).expect_err("writer byte budget above hard cap must fail");
ProxyConfig::load(&path).expect_err("writer byte budget above hard cap must fail");
let msg = err.to_string(); let msg = err.to_string();
assert!( assert!(
msg.contains("general.me_writer_byte_budget_bytes must be within [33570816, 268435456]"), msg.contains("general.me_writer_byte_budget_bytes must be within [33570816, 268435456]"),
@@ -165,8 +164,7 @@ direct_relay_buffer_budget_max_bytes = 16777217
"#, "#,
); );
let err = ProxyConfig::load(&path) let err = ProxyConfig::load(&path).expect_err("unaligned direct relay buffer budget must fail");
.expect_err("unaligned direct relay buffer budget must fail");
assert!( assert!(
err.to_string().contains( err.to_string().contains(
"general.direct_relay_buffer_budget_max_bytes must be 0 or a multiple of 4096" "general.direct_relay_buffer_budget_max_bytes must be 0 or a multiple of 4096"
@@ -184,8 +182,8 @@ direct_relay_buffer_budget_max_bytes = 2147487744
"#, "#,
); );
let err = ProxyConfig::load(&path) let err =
.expect_err("direct relay buffer budget above hard cap must fail"); ProxyConfig::load(&path).expect_err("direct relay buffer budget above hard cap must fail");
assert!(err.to_string().contains( assert!(err.to_string().contains(
"general.direct_relay_buffer_budget_max_bytes must be 0 or within [16777216, 2147483648]" "general.direct_relay_buffer_budget_max_bytes must be 0 or within [16777216, 2147483648]"
)); ));
+1 -2
View File
@@ -1126,8 +1126,7 @@ impl Default for GeneralConfig {
default_me_d2c_frame_buf_shrink_threshold_bytes(), default_me_d2c_frame_buf_shrink_threshold_bytes(),
direct_relay_copy_buf_c2s_bytes: default_direct_relay_copy_buf_c2s_bytes(), direct_relay_copy_buf_c2s_bytes: default_direct_relay_copy_buf_c2s_bytes(),
direct_relay_copy_buf_s2c_bytes: default_direct_relay_copy_buf_s2c_bytes(), direct_relay_copy_buf_s2c_bytes: default_direct_relay_copy_buf_s2c_bytes(),
direct_relay_buffer_budget_max_bytes: direct_relay_buffer_budget_max_bytes: default_direct_relay_buffer_budget_max_bytes(),
default_direct_relay_buffer_budget_max_bytes(),
me_warmup_stagger_enabled: default_true(), me_warmup_stagger_enabled: default_true(),
me_warmup_step_delay_ms: default_warmup_step_delay_ms(), me_warmup_step_delay_ms: default_warmup_step_delay_ms(),
me_warmup_step_jitter_ms: default_warmup_step_jitter_ms(), me_warmup_step_jitter_ms: default_warmup_step_jitter_ms(),
+4 -7
View File
@@ -33,11 +33,10 @@ use crate::conntrack_control;
use crate::crypto::SecureRandom; use crate::crypto::SecureRandom;
use crate::ip_tracker::UserIpTracker; use crate::ip_tracker::UserIpTracker;
use crate::network::probe::{decide_network_capabilities, log_probe_result, run_probe}; use crate::network::probe::{decide_network_capabilities, log_probe_result, run_probe};
use crate::proxy::route_mode::{RelayRouteMode, RouteRuntimeController};
use crate::proxy::direct_buffer_budget::{ use crate::proxy::direct_buffer_budget::{
DirectBufferBudget, resolve_direct_buffer_hard_limit, DirectBufferBudget, resolve_direct_buffer_hard_limit, spawn_direct_buffer_budget_controller,
spawn_direct_buffer_budget_controller,
}; };
use crate::proxy::route_mode::{RelayRouteMode, RouteRuntimeController};
use crate::proxy::shared_state::ProxySharedState; use crate::proxy::shared_state::ProxySharedState;
use crate::startup::{ use crate::startup::{
COMPONENT_API_BOOTSTRAP, COMPONENT_CONFIG_LOAD, COMPONENT_DC_CONNECTIVITY_PING, COMPONENT_API_BOOTSTRAP, COMPONENT_CONFIG_LOAD, COMPONENT_DC_CONNECTIVITY_PING,
@@ -477,10 +476,8 @@ async fn run_telemt_core(
config.network.dns_overrides.len() config.network.dns_overrides.len()
); );
} }
let direct_buffer_hard_limit = resolve_direct_buffer_hard_limit( let direct_buffer_hard_limit =
config.general.direct_relay_buffer_budget_max_bytes, resolve_direct_buffer_hard_limit(config.general.direct_relay_buffer_budget_max_bytes).await;
)
.await;
let direct_buffer_budget = DirectBufferBudget::new(direct_buffer_hard_limit); let direct_buffer_budget = DirectBufferBudget::new(direct_buffer_hard_limit);
info!( info!(
hard_limit_bytes = direct_buffer_hard_limit, hard_limit_bytes = direct_buffer_hard_limit,
+2 -8
View File
@@ -611,10 +611,7 @@ async fn render_metrics(
out, out,
"# HELP telemt_direct_relay_buffer_budget_bytes Direct relay copy-buffer budget and memory inputs" "# HELP telemt_direct_relay_buffer_budget_bytes Direct relay copy-buffer budget and memory inputs"
); );
let _ = writeln!( let _ = writeln!(out, "# TYPE telemt_direct_relay_buffer_budget_bytes gauge");
out,
"# TYPE telemt_direct_relay_buffer_budget_bytes gauge"
);
for (kind, value) in [ for (kind, value) in [
("hard_limit", direct_budget.hard_limit_bytes), ("hard_limit", direct_budget.hard_limit_bytes),
("target", direct_budget.target_bytes), ("target", direct_budget.target_bytes),
@@ -2517,10 +2514,7 @@ async fn render_metrics(
out, out,
"# HELP telemt_me_writer_byte_budget_limit_bytes Configured resident-memory budget per ME writer" "# HELP telemt_me_writer_byte_budget_limit_bytes Configured resident-memory budget per ME writer"
); );
let _ = writeln!( let _ = writeln!(out, "# TYPE telemt_me_writer_byte_budget_limit_bytes gauge");
out,
"# TYPE telemt_me_writer_byte_budget_limit_bytes gauge"
);
let _ = writeln!( let _ = writeln!(
out, out,
"telemt_me_writer_byte_budget_limit_bytes {}", "telemt_me_writer_byte_budget_limit_bytes {}",
+14 -21
View File
@@ -170,11 +170,9 @@ impl DirectBufferBudget {
} }
fn set_target_bytes(&self, target: u64) { fn set_target_bytes(&self, target: u64) {
let target = align_down( let target =
target align_down(target.clamp(self.target_floor_bytes(), self.hard_limit_bytes) as usize)
.clamp(self.target_floor_bytes(), self.hard_limit_bytes) as u64;
as usize,
) as u64;
let previous = self.target_bytes.swap(target, Ordering::AcqRel); let previous = self.target_bytes.swap(target, Ordering::AcqRel);
if target < previous { if target < previous {
let generation = self let generation = self
@@ -196,8 +194,7 @@ impl DirectBufferBudget {
/// Records a session that had to bypass the adaptive target at minimum size. /// Records a session that had to bypass the adaptive target at minimum size.
pub(crate) fn increment_minimum_fallback(&self) { pub(crate) fn increment_minimum_fallback(&self) {
self.minimum_fallback_total self.minimum_fallback_total.fetch_add(1, Ordering::Relaxed);
.fetch_add(1, Ordering::Relaxed);
} }
/// Records a session rejected because the absolute ceiling was exhausted. /// Records a session rejected because the absolute ceiling was exhausted.
@@ -369,7 +366,9 @@ pub(crate) fn spawn_direct_buffer_budget_controller(
budget.update_system_sample(sample); budget.update_system_sample(sample);
let snapshot = budget.snapshot(); let snapshot = budget.snapshot();
let denied_delta = snapshot.promotion_denied_total.saturating_sub(previous_denied); let denied_delta = snapshot
.promotion_denied_total
.saturating_sub(previous_denied);
previous_denied = snapshot.promotion_denied_total; previous_denied = snapshot.promotion_denied_total;
let fallback_delta = snapshot let fallback_delta = snapshot
.minimum_fallback_total .minimum_fallback_total
@@ -404,9 +403,7 @@ pub(crate) fn spawn_direct_buffer_budget_controller(
pool_snapshot.allocated, pool_snapshot.allocated,
pool_snapshot.allocated.saturating_sub(pool_snapshot.pooled), pool_snapshot.allocated.saturating_sub(pool_snapshot.pooled),
); );
stats.set_buffer_pool_replaced_nonstandard_total( stats.set_buffer_pool_replaced_nonstandard_total(pool_snapshot.replaced_nonstandard);
pool_snapshot.replaced_nonstandard,
);
let headroom_target = if sample.total_bytes == 0 { let headroom_target = if sample.total_bytes == 0 {
snapshot.hard_limit_bytes snapshot.hard_limit_bytes
@@ -454,11 +451,8 @@ fn connection_fill_pct(stats: &Stats, max_connections: u32) -> Option<u8> {
return None; return None;
} }
Some( Some(
((stats ((stats.get_current_connections_total().saturating_mul(100)) / u64::from(max_connections))
.get_current_connections_total() .min(100) as u8,
.saturating_mul(100))
/ u64::from(max_connections))
.min(100) as u8,
) )
} }
@@ -485,8 +479,7 @@ async fn read_system_memory_sample() -> SystemMemorySample {
let cgroup_v2_max = read_cgroup_limit("/sys/fs/cgroup/memory.max").await; let cgroup_v2_max = read_cgroup_limit("/sys/fs/cgroup/memory.max").await;
let cgroup_v2_current = read_u64_file("/sys/fs/cgroup/memory.current").await; let cgroup_v2_current = read_u64_file("/sys/fs/cgroup/memory.current").await;
let cgroup_v1_max = read_cgroup_limit("/sys/fs/cgroup/memory/memory.limit_in_bytes").await; let cgroup_v1_max = read_cgroup_limit("/sys/fs/cgroup/memory/memory.limit_in_bytes").await;
let cgroup_v1_current = let cgroup_v1_current = read_u64_file("/sys/fs/cgroup/memory/memory.usage_in_bytes").await;
read_u64_file("/sys/fs/cgroup/memory/memory.usage_in_bytes").await;
let cgroup_max = cgroup_v2_max.or(cgroup_v1_max); let cgroup_max = cgroup_v2_max.or(cgroup_v1_max);
let cgroup_current = cgroup_v2_current.or(cgroup_v1_current); let cgroup_current = cgroup_v2_current.or(cgroup_v1_current);
@@ -495,9 +488,9 @@ async fn read_system_memory_sample() -> SystemMemorySample {
(host, Some(limit)) => host.min(limit), (host, Some(limit)) => host.min(limit),
(host, None) => host, (host, None) => host,
}; };
let cgroup_available = cgroup_max.zip(cgroup_current).map(|(limit, current)| { let cgroup_available = cgroup_max
limit.saturating_sub(current) .zip(cgroup_current)
}); .map(|(limit, current)| limit.saturating_sub(current));
let available = match (host_available, cgroup_available) { let available = match (host_available, cgroup_available) {
(0, Some(value)) => value, (0, Some(value)) => value,
(host, Some(value)) => host.min(value), (host, Some(value)) => host.min(value),
+15 -15
View File
@@ -346,21 +346,21 @@ where
Duration::from_secs(1800) Duration::from_secs(1800)
}; };
let relay_result = crate::proxy::relay::relay_direct_adaptive( let relay_result = crate::proxy::relay::relay_direct_adaptive(
client_reader, client_reader,
client_writer, client_writer,
tg_reader, tg_reader,
tg_writer, tg_writer,
config.general.direct_relay_copy_buf_c2s_bytes, config.general.direct_relay_copy_buf_c2s_bytes,
config.general.direct_relay_copy_buf_s2c_bytes, config.general.direct_relay_copy_buf_s2c_bytes,
config.server.max_connections, config.server.max_connections,
user, user,
Arc::clone(&stats), Arc::clone(&stats),
config.access.user_data_quota.get(user).copied(), config.access.user_data_quota.get(user).copied(),
traffic_lease, traffic_lease,
relay_activity_timeout, relay_activity_timeout,
session_cancel.clone(), session_cancel.clone(),
Arc::clone(&shared.direct_buffer_budget), Arc::clone(&shared.direct_buffer_budget),
); );
tokio::pin!(relay_result); tokio::pin!(relay_result);
let relay_result = loop { let relay_result = loop {
if let Some(cutover) = if let Some(cutover) =
+1 -1
View File
@@ -84,8 +84,8 @@ fn watchdog_delta(current: u64, previous: u64) -> u64 {
current.saturating_sub(previous) current.saturating_sub(previous)
} }
mod io;
mod adaptive_copy; mod adaptive_copy;
mod io;
pub(crate) use self::adaptive_copy::relay_direct_adaptive; pub(crate) use self::adaptive_copy::relay_direct_adaptive;
+2 -7
View File
@@ -488,13 +488,8 @@ fn apply_global_pressure_demotion(
return; return;
} }
*controller = SessionAdaptiveController::new(target); *controller = SessionAdaptiveController::new(target);
let sizes = direct_copy_buffers_for_tier_with_ceilings( let sizes =
target, direct_copy_buffers_for_tier_with_ceilings(target, base.0, base.1, ceilings.0, ceilings.1);
base.0,
base.1,
ceilings.0,
ceilings.1,
);
set_desired_sizes(c2s_state, s2c_state, sizes); set_desired_sizes(c2s_state, s2c_state, sizes);
lease.set_tier(target.as_u8() as usize); lease.set_tier(target.as_u8() as usize);
budget.increment_global_pressure_demotion(); budget.increment_global_pressure_demotion();
+1 -3
View File
@@ -9,10 +9,8 @@ use dashmap::DashMap;
use tokio::sync::{OwnedSemaphorePermit, Semaphore, mpsc}; use tokio::sync::{OwnedSemaphorePermit, Semaphore, mpsc};
use tokio_util::sync::CancellationToken; use tokio_util::sync::CancellationToken;
use crate::proxy::direct_buffer_budget::{DirectBufferBudget, fallback_direct_buffer_hard_limit};
use crate::proxy::handshake::{AuthProbeSaturationState, AuthProbeState}; use crate::proxy::handshake::{AuthProbeSaturationState, AuthProbeState};
use crate::proxy::direct_buffer_budget::{
DirectBufferBudget, fallback_direct_buffer_hard_limit,
};
use crate::proxy::middle_relay::{DesyncDedupRotationState, RelayIdleCandidateRegistry}; use crate::proxy::middle_relay::{DesyncDedupRotationState, RelayIdleCandidateRegistry};
use crate::proxy::traffic_limiter::TrafficLimiter; use crate::proxy::traffic_limiter::TrafficLimiter;
+2 -1
View File
@@ -292,7 +292,8 @@ impl Stats {
} }
/// Returns the count of blocking writer byte-budget waits. /// Returns the count of blocking writer byte-budget waits.
pub fn get_me_writer_byte_budget_wait_total(&self) -> u64 { pub fn get_me_writer_byte_budget_wait_total(&self) -> u64 {
self.me_writer_byte_budget_wait_total.load(Ordering::Relaxed) self.me_writer_byte_budget_wait_total
.load(Ordering::Relaxed)
} }
/// Returns the count of writer byte-budget wait timeouts. /// Returns the count of writer byte-budget wait timeouts.
pub fn get_me_writer_byte_budget_timeout_total(&self) -> u64 { pub fn get_me_writer_byte_budget_timeout_total(&self) -> u64 {
+5 -5
View File
@@ -114,11 +114,11 @@ impl Stats {
); );
} }
pub(crate) fn release_me_writer_byte_budget_inflight_bytes(&self, bytes: u64) { pub(crate) fn release_me_writer_byte_budget_inflight_bytes(&self, bytes: u64) {
let _ = self.me_writer_byte_budget_inflight_bytes_gauge.fetch_update( let _ = self
Ordering::Relaxed, .me_writer_byte_budget_inflight_bytes_gauge
Ordering::Relaxed, .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
|current| Some(current.saturating_sub(bytes)), Some(current.saturating_sub(bytes))
); });
} }
pub(crate) fn increment_me_writer_byte_budget_wait_total(&self) { pub(crate) fn increment_me_writer_byte_budget_wait_total(&self) {
if self.telemetry_me_allows_normal() { if self.telemetry_me_allows_normal() {
+4 -1
View File
@@ -400,7 +400,10 @@ mod tests {
.expect("writer byte permits must be available"); .expect("writer byte permits must be available");
let mut writer_permit = WriterBytePermit::new(permit, 32 * 1024, stats.clone()); let mut writer_permit = WriterBytePermit::new(permit, 32 * 1024, stats.clone());
assert_eq!(stats.get_me_writer_byte_budget_queued_bytes_gauge(), 32 * 1024); assert_eq!(
stats.get_me_writer_byte_budget_queued_bytes_gauge(),
32 * 1024
);
assert_eq!(stats.get_me_writer_byte_budget_inflight_bytes_gauge(), 0); assert_eq!(stats.get_me_writer_byte_budget_inflight_bytes_gauge(), 0);
writer_permit.mark_inflight(); writer_permit.mark_inflight();
@@ -77,13 +77,9 @@ impl ConnRegistry {
.writer_idle_since_epoch_secs .writer_idle_since_epoch_secs
.entry(writer_id) .entry(writer_id)
.or_insert_with(Self::now_epoch_secs); .or_insert_with(Self::now_epoch_secs);
self.writers.map.insert( self.writers
writer_id, .map
super::WriterRoute { .insert(writer_id, super::WriterRoute { tx, byte_budget });
tx,
byte_budget,
},
);
} }
/// Unregister connection, returning associated writer_id if any. /// Unregister connection, returning associated writer_id if any.
+47 -67
View File
@@ -6,8 +6,8 @@ use std::sync::Arc;
use std::sync::atomic::Ordering; use std::sync::atomic::Ordering;
use std::time::{Duration, Instant}; use std::time::{Duration, Instant};
use tokio::sync::{OwnedSemaphorePermit, Semaphore, TryAcquireError, mpsc};
use tokio::sync::mpsc::error::TrySendError; use tokio::sync::mpsc::error::TrySendError;
use tokio::sync::{OwnedSemaphorePermit, Semaphore, TryAcquireError, mpsc};
use tracing::{debug, warn}; use tracing::{debug, warn};
use super::MePool; use super::MePool;
@@ -74,16 +74,15 @@ async fn reserve_writer_command_slot(
) -> std::result::Result<mpsc::OwnedPermit<WriterCommand>, WriterCommandReserveError> { ) -> std::result::Result<mpsc::OwnedPermit<WriterCommand>, WriterCommandReserveError> {
let reserve = tx.clone().reserve_owned(); let reserve = tx.clone().reserve_owned();
match deadline { match deadline {
Some(deadline) => match tokio::time::timeout( Some(deadline) => {
deadline.saturating_duration_since(Instant::now()), match tokio::time::timeout(deadline.saturating_duration_since(Instant::now()), reserve)
reserve, .await
) {
.await Ok(Ok(permit)) => Ok(permit),
{ Ok(Err(_)) => Err(WriterCommandReserveError::Closed),
Ok(Ok(permit)) => Ok(permit), Err(_) => Err(WriterCommandReserveError::TimedOut),
Ok(Err(_)) => Err(WriterCommandReserveError::Closed), }
Err(_) => Err(WriterCommandReserveError::TimedOut), }
},
None => reserve.await.map_err(|_| WriterCommandReserveError::Closed), None => reserve.await.map_err(|_| WriterCommandReserveError::Closed),
} }
} }
@@ -102,7 +101,10 @@ fn writer_resident_permits(
let permits = resident_bytes.div_ceil(ME_WRITER_BYTE_PERMIT_UNIT_BYTES); let permits = resident_bytes.div_ceil(ME_WRITER_BYTE_PERMIT_UNIT_BYTES);
let permits = u32::try_from(permits).ok()?; let permits = u32::try_from(permits).ok()?;
let reserved_bytes = (permits as usize).checked_mul(ME_WRITER_BYTE_PERMIT_UNIT_BYTES)?; let reserved_bytes = (permits as usize).checked_mul(ME_WRITER_BYTE_PERMIT_UNIT_BYTES)?;
Some((permits.max(1), reserved_bytes.max(ME_WRITER_BYTE_PERMIT_UNIT_BYTES))) Some((
permits.max(1),
reserved_bytes.max(ME_WRITER_BYTE_PERMIT_UNIT_BYTES),
))
} }
fn proxy_req_resident_permits( fn proxy_req_resident_permits(
@@ -146,23 +148,18 @@ async fn reserve_writer_bytes(
let acquire = byte_budget.clone().acquire_many_owned(permits); let acquire = byte_budget.clone().acquire_many_owned(permits);
match deadline { match deadline {
Some(deadline) => match tokio::time::timeout( Some(deadline) => {
deadline.saturating_duration_since(Instant::now()), match tokio::time::timeout(deadline.saturating_duration_since(Instant::now()), acquire)
acquire, .await
) {
.await Ok(Ok(permit)) => Ok(WriterBytePermit::new(permit, reserved_bytes, stats.clone())),
{ Ok(Err(_)) => Err(WriterByteReserveError::Closed),
Ok(Ok(permit)) => Ok(WriterBytePermit::new( Err(_) => {
permit, stats.increment_me_writer_byte_budget_timeout_total();
reserved_bytes, Err(WriterByteReserveError::TimedOut)
stats.clone(), }
)),
Ok(Err(_)) => Err(WriterByteReserveError::Closed),
Err(_) => {
stats.increment_me_writer_byte_budget_timeout_total();
Err(WriterByteReserveError::TimedOut)
} }
}, }
None => acquire None => acquire
.await .await
.map(|permit| WriterBytePermit::new(permit, reserved_bytes, stats.clone())) .map(|permit| WriterBytePermit::new(permit, reserved_bytes, stats.clone()))
@@ -189,27 +186,21 @@ impl MePool {
.len() .len()
.checked_add(LEGACY_PROXY_REQ_SOURCE_CAPACITY_OVERHEAD_BYTES) .checked_add(LEGACY_PROXY_REQ_SOURCE_CAPACITY_OVERHEAD_BYTES)
else { else {
self.stats self.stats.increment_me_writer_byte_budget_oversize_total();
.increment_me_writer_byte_budget_oversize_total();
return Err(ProxyError::Proxy( return Err(ProxyError::Proxy(
"ME writer payload residency calculation overflow".into(), "ME writer payload residency calculation overflow".into(),
)); ));
}; };
let Some((writer_byte_permits, writer_reserved_bytes)) = proxy_req_resident_permits( let Some((writer_byte_permits, writer_reserved_bytes)) =
source_capacity, proxy_req_resident_permits(source_capacity, data.len(), tag, proto_flags)
data.len(), else {
tag, self.stats.increment_me_writer_byte_budget_oversize_total();
proto_flags,
) else {
self.stats
.increment_me_writer_byte_budget_oversize_total();
return Err(ProxyError::Proxy( return Err(ProxyError::Proxy(
"ME writer payload residency calculation overflow".into(), "ME writer payload residency calculation overflow".into(),
)); ));
}; };
if writer_byte_permits as usize > self.writer_lifecycle.writer_byte_budget_permits { if writer_byte_permits as usize > self.writer_lifecycle.writer_byte_budget_permits {
self.stats self.stats.increment_me_writer_byte_budget_oversize_total();
.increment_me_writer_byte_budget_oversize_total();
return Err(ProxyError::Proxy( return Err(ProxyError::Proxy(
"ME writer payload exceeds configured byte budget".into(), "ME writer payload exceeds configured byte budget".into(),
)); ));
@@ -254,9 +245,8 @@ impl MePool {
loop { loop {
if let Some((current, current_meta)) = self.registry.get_writer_with_meta(conn_id).await if let Some((current, current_meta)) = self.registry.get_writer_with_meta(conn_id).await
{ {
let deadline = writer_send_deadline( let deadline =
self.route_runtime.me_route_blocking_send_timeout, writer_send_deadline(self.route_runtime.me_route_blocking_send_timeout);
);
let writer_permit = match reserve_writer_bytes( let writer_permit = match reserve_writer_bytes(
&current.byte_budget, &current.byte_budget,
writer_byte_permits, writer_byte_permits,
@@ -275,7 +265,10 @@ impl MePool {
)); ));
} }
Err(WriterByteReserveError::Closed) => { Err(WriterByteReserveError::Closed) => {
warn!(writer_id = current.writer_id, "ME writer byte budget closed"); warn!(
writer_id = current.writer_id,
"ME writer byte budget closed"
);
self.remove_writer_and_close_clients(current.writer_id) self.remove_writer_and_close_clients(current.writer_id)
.await; .await;
continue; continue;
@@ -293,12 +286,7 @@ impl MePool {
return Ok(()); return Ok(());
} }
Err(TrySendError::Full(cmd)) => { Err(TrySendError::Full(cmd)) => {
match reserve_writer_command_slot( match reserve_writer_command_slot(&current.tx, deadline).await {
&current.tx,
deadline,
)
.await
{
Ok(permit) => { Ok(permit) => {
permit.send(cmd); permit.send(cmd);
self.note_hybrid_route_success(); self.note_hybrid_route_success();
@@ -716,9 +704,7 @@ impl MePool {
} }
self.stats self.stats
.increment_me_writer_pick_blocking_fallback_total(); .increment_me_writer_pick_blocking_fallback_total();
let deadline = writer_send_deadline( let deadline = writer_send_deadline(self.route_runtime.me_route_blocking_send_timeout);
self.route_runtime.me_route_blocking_send_timeout,
);
let writer_permit = match reserve_writer_bytes( let writer_permit = match reserve_writer_bytes(
&w.byte_budget, &w.byte_budget,
writer_byte_permits, writer_byte_permits,
@@ -801,24 +787,20 @@ impl MePool {
tag.as_ref().map(|tag| tag.as_slice()), tag.as_ref().map(|tag| tag.as_slice()),
proto_flags, proto_flags,
) else { ) else {
self.stats self.stats.increment_me_writer_byte_budget_oversize_total();
.increment_me_writer_byte_budget_oversize_total();
return Err(ProxyError::Proxy( return Err(ProxyError::Proxy(
"ME writer payload residency calculation overflow".into(), "ME writer payload residency calculation overflow".into(),
)); ));
}; };
if writer_byte_permits as usize > self.writer_lifecycle.writer_byte_budget_permits { if writer_byte_permits as usize > self.writer_lifecycle.writer_byte_budget_permits {
self.stats self.stats.increment_me_writer_byte_budget_oversize_total();
.increment_me_writer_byte_budget_oversize_total();
return Err(ProxyError::Proxy( return Err(ProxyError::Proxy(
"ME writer payload exceeds configured byte budget".into(), "ME writer payload exceeds configured byte budget".into(),
)); ));
} }
if let Some((current, current_meta)) = self.registry.get_writer_with_meta(conn_id).await { if let Some((current, current_meta)) = self.registry.get_writer_with_meta(conn_id).await {
let deadline = writer_send_deadline( let deadline = writer_send_deadline(self.route_runtime.me_route_blocking_send_timeout);
self.route_runtime.me_route_blocking_send_timeout,
);
let writer_permit = match reserve_writer_bytes( let writer_permit = match reserve_writer_bytes(
&current.byte_budget, &current.byte_budget,
writer_byte_permits, writer_byte_permits,
@@ -837,7 +819,10 @@ impl MePool {
)); ));
} }
Err(WriterByteReserveError::Closed) => { Err(WriterByteReserveError::Closed) => {
warn!(writer_id = current.writer_id, "ME writer byte budget closed"); warn!(
writer_id = current.writer_id,
"ME writer byte budget closed"
);
self.remove_writer_and_close_clients(current.writer_id) self.remove_writer_and_close_clients(current.writer_id)
.await; .await;
return self return self
@@ -870,12 +855,7 @@ impl MePool {
return Ok(()); return Ok(());
} }
Err(TrySendError::Full(cmd)) => { Err(TrySendError::Full(cmd)) => {
match reserve_writer_command_slot( match reserve_writer_command_slot(&current.tx, deadline).await {
&current.tx,
deadline,
)
.await
{
Ok(permit) => { Ok(permit) => {
permit.send(cmd); permit.send(cmd);
self.note_hybrid_route_success(); self.note_hybrid_route_success();