mirror of https://github.com/ospab/ostp.git
Compare commits
No commits in common. "cd12b01bc356d7b95ceb09811da875f373a40097" and "66a1e9784039ad989dae2519ce078e36137f13c7" have entirely different histories.
cd12b01bc3
...
66a1e97840
|
|
@ -2,5 +2,5 @@
|
||||||
"target_version": "0.4.2",
|
"target_version": "0.4.2",
|
||||||
"branch": "beta",
|
"branch": "beta",
|
||||||
"alpha_iteration": 0,
|
"alpha_iteration": 0,
|
||||||
"beta_iteration": 2
|
"beta_iteration": 1
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -16,7 +16,7 @@ publish_to: 'none' # Remove this line if you wish to publish to pub.dev
|
||||||
# https://developer.apple.com/library/archive/documentation/General/Reference/InfoPlistKeyReference/Articles/CoreFoundationKeys.html
|
# https://developer.apple.com/library/archive/documentation/General/Reference/InfoPlistKeyReference/Articles/CoreFoundationKeys.html
|
||||||
# In Windows, build-name is used as the major, minor, and patch parts
|
# In Windows, build-name is used as the major, minor, and patch parts
|
||||||
# of the product and file versions while build-number is used as the build suffix.
|
# of the product and file versions while build-number is used as the build suffix.
|
||||||
version: 0.4.2+21
|
version: 0.4.2+20
|
||||||
|
|
||||||
environment:
|
environment:
|
||||||
sdk: ^3.11.4
|
sdk: ^3.11.4
|
||||||
|
|
|
||||||
|
|
@ -402,6 +402,55 @@
|
||||||
</select>
|
</select>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
|
<!-- Advanced TCP/UoT Settings (visible only if uot is selected) -->
|
||||||
|
<div id="pm-tcp-settings" style="display:none; padding: 10px; background: rgba(0,0,0,0.2); border-radius: 8px; margin-bottom: 15px;">
|
||||||
|
<div class="toggle-row" style="padding:0; border:none; margin-bottom:10px;">
|
||||||
|
<div class="toggle-text">
|
||||||
|
<span class="toggle-name">TCP Fragmentation</span>
|
||||||
|
<span class="toggle-hint">Split handshake to bypass DPI</span>
|
||||||
|
</div>
|
||||||
|
<label class="toggle">
|
||||||
|
<input type="checkbox" id="pm-tcp-frag" />
|
||||||
|
<span class="toggle-track"><span class="toggle-thumb"></span></span>
|
||||||
|
</label>
|
||||||
|
</div>
|
||||||
|
|
||||||
|
<div id="pm-frag-details" style="display:none;">
|
||||||
|
<div style="display:flex; gap:10px; margin-bottom:10px;">
|
||||||
|
<div class="inline-field" style="padding:0; border:none; flex:1;">
|
||||||
|
<span class="field-label">Chunk Size</span>
|
||||||
|
<input id="pm-frag-chunk" class="field-input compact" type="number" placeholder="2" min="1" />
|
||||||
|
</div>
|
||||||
|
<div class="inline-field" style="padding:0; border:none; flex:1;">
|
||||||
|
<span class="field-label">Sleep (ms)</span>
|
||||||
|
<input id="pm-frag-sleep" class="field-input compact" type="number" placeholder="2" min="0" />
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
|
||||||
|
<div class="section-divider-mini" style="margin-top:0;"><span>Junk Packets</span></div>
|
||||||
|
<div style="display:flex; gap:10px; margin-bottom:10px;">
|
||||||
|
<div class="inline-field" style="padding:0; border:none; flex:1;">
|
||||||
|
<span class="field-label">Count (Min)</span>
|
||||||
|
<input id="pm-junk-pc-min" class="field-input compact" type="number" placeholder="2" min="0" />
|
||||||
|
</div>
|
||||||
|
<div class="inline-field" style="padding:0; border:none; flex:1;">
|
||||||
|
<span class="field-label">Count (Max)</span>
|
||||||
|
<input id="pm-junk-pc-max" class="field-input compact" type="number" placeholder="5" min="0" />
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
<div style="display:flex; gap:10px;">
|
||||||
|
<div class="inline-field" style="padding:0; border:none; flex:1;">
|
||||||
|
<span class="field-label">Size (Min)</span>
|
||||||
|
<input id="pm-junk-ps-min" class="field-input compact" type="number" placeholder="100" min="0" />
|
||||||
|
</div>
|
||||||
|
<div class="inline-field" style="padding:0; border:none; flex:1;">
|
||||||
|
<span class="field-label">Size (Max)</span>
|
||||||
|
<input id="pm-junk-ps-max" class="field-input compact" type="number" placeholder="1000" min="0" />
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
|
||||||
<div class="modal-actions">
|
<div class="modal-actions">
|
||||||
<button id="btn-profile-cancel" class="btn secondary">Cancel</button>
|
<button id="btn-profile-cancel" class="btn secondary">Cancel</button>
|
||||||
<button id="btn-profile-delete" class="btn danger" style="display:none;">Delete</button>
|
<button id="btn-profile-delete" class="btn danger" style="display:none;">Delete</button>
|
||||||
|
|
|
||||||
|
|
@ -113,6 +113,15 @@ const pmName = $('pm-name');
|
||||||
const pmServer = $('pm-server');
|
const pmServer = $('pm-server');
|
||||||
const pmKey = $('pm-key');
|
const pmKey = $('pm-key');
|
||||||
const pmTransport = $('pm-transport');
|
const pmTransport = $('pm-transport');
|
||||||
|
const pmTcpFrag = $('pm-tcp-frag');
|
||||||
|
const pmFragChunk = $('pm-frag-chunk');
|
||||||
|
const pmFragSleep = $('pm-frag-sleep');
|
||||||
|
const pmJunkPcMin = $('pm-junk-pc-min');
|
||||||
|
const pmJunkPcMax = $('pm-junk-pc-max');
|
||||||
|
const pmJunkPsMin = $('pm-junk-ps-min');
|
||||||
|
const pmJunkPsMax = $('pm-junk-ps-max');
|
||||||
|
const pmTcpSettings = $('pm-tcp-settings');
|
||||||
|
const pmFragDetails = $('pm-frag-details');
|
||||||
const btnProfileCancel = $('btn-profile-cancel');
|
const btnProfileCancel = $('btn-profile-cancel');
|
||||||
const btnProfileSave = $('btn-profile-save');
|
const btnProfileSave = $('btn-profile-save');
|
||||||
const btnProfileDelete = $('btn-profile-delete');
|
const btnProfileDelete = $('btn-profile-delete');
|
||||||
|
|
@ -513,15 +522,31 @@ function openProfileEditor(id) {
|
||||||
pmServer.value = p.server || '';
|
pmServer.value = p.server || '';
|
||||||
pmKey.value = p.key || '';
|
pmKey.value = p.key || '';
|
||||||
pmTransport.value = p.transport || 'udp';
|
pmTransport.value = p.transport || 'udp';
|
||||||
|
pmTcpFrag.checked = !!p.tcp_fragmentation;
|
||||||
|
pmFragChunk.value = p.frag_chunk || 2;
|
||||||
|
pmFragSleep.value = p.frag_sleep || 2;
|
||||||
|
pmJunkPcMin.value = p.junk_pc ? p.junk_pc[0] : 2;
|
||||||
|
pmJunkPcMax.value = p.junk_pc ? p.junk_pc[1] : 5;
|
||||||
|
pmJunkPsMin.value = p.junk_ps ? p.junk_ps[0] : 100;
|
||||||
|
pmJunkPsMax.value = p.junk_ps ? p.junk_ps[1] : 1000;
|
||||||
btnProfileDelete.style.display = '';
|
btnProfileDelete.style.display = '';
|
||||||
} else {
|
} else {
|
||||||
profileModalTitle.textContent = 'New Profile';
|
profileModalTitle.textContent = 'New Profile';
|
||||||
pmName.value = pmServer.value = pmKey.value = '';
|
pmName.value = pmServer.value = pmKey.value = '';
|
||||||
pmTransport.value = 'udp';
|
pmTransport.value = 'udp';
|
||||||
|
pmTcpFrag.checked = false;
|
||||||
|
pmFragChunk.value = 2;
|
||||||
|
pmFragSleep.value = 2;
|
||||||
|
pmJunkPcMin.value = 2;
|
||||||
|
pmJunkPcMax.value = 5;
|
||||||
|
pmJunkPsMin.value = 100;
|
||||||
|
pmJunkPsMax.value = 1000;
|
||||||
btnProfileDelete.style.display = 'none';
|
btnProfileDelete.style.display = 'none';
|
||||||
}
|
}
|
||||||
pmKey.type = 'password';
|
pmKey.type = 'password';
|
||||||
profileModal.classList.remove('hidden');
|
profileModal.classList.remove('hidden');
|
||||||
|
pmTransport.dispatchEvent(new Event('change'));
|
||||||
|
pmTcpFrag.dispatchEvent(new Event('change'));
|
||||||
setTimeout(() => pmName.focus(), 80);
|
setTimeout(() => pmName.focus(), 80);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -539,6 +564,11 @@ function saveProfileFromEditor() {
|
||||||
server,
|
server,
|
||||||
key,
|
key,
|
||||||
transport: pmTransport.value,
|
transport: pmTransport.value,
|
||||||
|
tcp_fragmentation: pmTcpFrag.checked,
|
||||||
|
frag_chunk: parseInt(pmFragChunk.value) || 2,
|
||||||
|
frag_sleep: parseInt(pmFragSleep.value) || 2,
|
||||||
|
junk_pc: [parseInt(pmJunkPcMin.value)||2, parseInt(pmJunkPcMax.value)||5],
|
||||||
|
junk_ps: [parseInt(pmJunkPsMin.value)||100, parseInt(pmJunkPsMax.value)||1000],
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
|
|
@ -548,6 +578,11 @@ function saveProfileFromEditor() {
|
||||||
server,
|
server,
|
||||||
key,
|
key,
|
||||||
transport: pmTransport.value,
|
transport: pmTransport.value,
|
||||||
|
tcp_fragmentation: pmTcpFrag.checked,
|
||||||
|
frag_chunk: parseInt(pmFragChunk.value) || 2,
|
||||||
|
frag_sleep: parseInt(pmFragSleep.value) || 2,
|
||||||
|
junk_pc: [parseInt(pmJunkPcMin.value)||2, parseInt(pmJunkPcMax.value)||5],
|
||||||
|
junk_ps: [parseInt(pmJunkPsMin.value)||100, parseInt(pmJunkPsMax.value)||1000],
|
||||||
};
|
};
|
||||||
profiles.push(p);
|
profiles.push(p);
|
||||||
if (!activeId) { activeId = p.id; saveActiveId(activeId); }
|
if (!activeId) { activeId = p.id; saveActiveId(activeId); }
|
||||||
|
|
@ -840,6 +875,12 @@ window.addEventListener('DOMContentLoaded', async () => {
|
||||||
btnProfileCancel.addEventListener('click', () => profileModal.classList.add('hidden'));
|
btnProfileCancel.addEventListener('click', () => profileModal.classList.add('hidden'));
|
||||||
btnProfileSave.addEventListener('click', saveProfileFromEditor);
|
btnProfileSave.addEventListener('click', saveProfileFromEditor);
|
||||||
btnProfileDelete.addEventListener('click', deleteEditingProfile);
|
btnProfileDelete.addEventListener('click', deleteEditingProfile);
|
||||||
|
pmTransport.addEventListener('change', () => {
|
||||||
|
pmTcpSettings.style.display = pmTransport.value === 'uot' ? 'block' : 'none';
|
||||||
|
});
|
||||||
|
pmTcpFrag.addEventListener('change', () => {
|
||||||
|
pmFragDetails.style.display = pmTcpFrag.checked ? 'block' : 'none';
|
||||||
|
});
|
||||||
btnPeekPm.addEventListener('click', () => {
|
btnPeekPm.addEventListener('click', () => {
|
||||||
pmKey.type = pmKey.type === 'password' ? 'text' : 'password';
|
pmKey.type = pmKey.type === 'password' ? 'text' : 'password';
|
||||||
});
|
});
|
||||||
|
|
|
||||||
|
|
@ -246,30 +246,6 @@ impl Dispatcher {
|
||||||
self.peer_machines.len()
|
self.peer_machines.len()
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Per-session download-direction congestion headroom, in packets:
|
|
||||||
/// `(session_id, available)` where `available = clamped cwnd - in_flight`.
|
|
||||||
///
|
|
||||||
/// Consumed by the relay's per-target-connection reader tasks (see
|
|
||||||
/// `relay::handle_relay_message`'s Connect handler) to throttle how fast
|
|
||||||
/// they pull bytes from the upstream target and forward them to the
|
|
||||||
/// client's OSTP session. Without this, a fast target (e.g. a CDN) gets
|
|
||||||
/// read and forwarded as fast as the target can serve, completely
|
|
||||||
/// ignoring the client-facing session's real congestion window - on a
|
|
||||||
/// lossy/jittery client path that self-inflicts a loss burst, which
|
|
||||||
/// wrecks the RTT/RTO estimate and can stall the session hard enough to
|
|
||||||
/// trip the client's keepalive reconnect. Same clamp(16, 16384) the
|
|
||||||
/// client uses for its own analogous uplink gate, for symmetry.
|
|
||||||
pub fn snapshot_backpressure(&self) -> Vec<(u32, i64)> {
|
|
||||||
self.peer_machines
|
|
||||||
.iter()
|
|
||||||
.map(|(&sid, ps)| {
|
|
||||||
let cwnd = (ps.machine.cwnd_packets() as i64).clamp(16, 16384);
|
|
||||||
let in_flight = ps.machine.in_flight_count() as i64;
|
|
||||||
(sid, cwnd - in_flight)
|
|
||||||
})
|
|
||||||
.collect()
|
|
||||||
}
|
|
||||||
|
|
||||||
pub fn on_datagram(&mut self, peer: SocketAddr, packet: Bytes) -> Result<DispatchOutcome> {
|
pub fn on_datagram(&mut self, peer: SocketAddr, packet: Bytes) -> Result<DispatchOutcome> {
|
||||||
if packet.len() < 4 {
|
if packet.len() < 4 {
|
||||||
return Ok(DispatchOutcome::Unauthorized);
|
return Ok(DispatchOutcome::Unauthorized);
|
||||||
|
|
|
||||||
|
|
@ -1,9 +1,7 @@
|
||||||
use anyhow::Result;
|
use anyhow::Result;
|
||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
use portable_atomic::AtomicI64;
|
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::net::IpAddr;
|
use std::net::IpAddr;
|
||||||
use std::sync::{Arc, RwLock};
|
|
||||||
|
|
||||||
use dispatcher::{DispatchOutcome, Dispatcher};
|
use dispatcher::{DispatchOutcome, Dispatcher};
|
||||||
use ostp_core::relay::RelayMessage;
|
use ostp_core::relay::RelayMessage;
|
||||||
|
|
@ -12,12 +10,6 @@ use tokio::net::UdpSocket;
|
||||||
use tokio::sync::mpsc;
|
use tokio::sync::mpsc;
|
||||||
use tokio::time::{interval, Duration, Instant};
|
use tokio::time::{interval, Duration, Instant};
|
||||||
|
|
||||||
/// Shared per-session download-direction congestion headroom (packets),
|
|
||||||
/// published by `handle_tick` from `Dispatcher::snapshot_backpressure` and
|
|
||||||
/// read lock-free by relay reader tasks. See that method's doc comment for
|
|
||||||
/// why this exists.
|
|
||||||
pub(crate) type SessionBackpressure = Arc<RwLock<HashMap<u32, Arc<AtomicI64>>>>;
|
|
||||||
|
|
||||||
mod dispatcher;
|
mod dispatcher;
|
||||||
pub mod outbound;
|
pub mod outbound;
|
||||||
pub mod api;
|
pub mod api;
|
||||||
|
|
@ -475,7 +467,6 @@ async fn run_server_loop(
|
||||||
let mut last_empty_app_log = Instant::now() - Duration::from_secs(10);
|
let mut last_empty_app_log = Instant::now() - Duration::from_secs(10);
|
||||||
let mut peer_last_seen: HashMap<IpAddr, Instant> = HashMap::new();
|
let mut peer_last_seen: HashMap<IpAddr, Instant> = HashMap::new();
|
||||||
let mut peer_available: HashMap<IpAddr, bool> = HashMap::new();
|
let mut peer_available: HashMap<IpAddr, bool> = HashMap::new();
|
||||||
let session_backpressure: SessionBackpressure = Arc::new(RwLock::new(HashMap::new()));
|
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
|
|
@ -498,8 +489,7 @@ async fn run_server_loop(
|
||||||
packet, peer, &mut dispatcher, &tcp_map, &socket, &mut remotes, &ui_event_tx,
|
packet, peer, &mut dispatcher, &tcp_map, &socket, &mut remotes, &ui_event_tx,
|
||||||
stream_tx.clone(), udp_reply_tx.clone(), connect_tx.clone(),
|
stream_tx.clone(), udp_reply_tx.clone(), connect_tx.clone(),
|
||||||
router.clone(),
|
router.clone(),
|
||||||
&mut peer_last_seen, &mut peer_available, &mut last_empty_app_log,
|
&mut peer_last_seen, &mut peer_available, &mut last_empty_app_log
|
||||||
&session_backpressure
|
|
||||||
).await {
|
).await {
|
||||||
tracing::error!("handle_udp_packet error: {}", e);
|
tracing::error!("handle_udp_packet error: {}", e);
|
||||||
}
|
}
|
||||||
|
|
@ -543,7 +533,7 @@ async fn run_server_loop(
|
||||||
_ = retransmit_tick.tick() => {
|
_ = retransmit_tick.tick() => {
|
||||||
if let Err(e) = handle_tick(
|
if let Err(e) = handle_tick(
|
||||||
&mut dispatcher, &tcp_map, &socket, &mut remotes, &ui_event_tx,
|
&mut dispatcher, &tcp_map, &socket, &mut remotes, &ui_event_tx,
|
||||||
&mut peer_last_seen, &mut peer_available, &session_backpressure
|
&mut peer_last_seen, &mut peer_available
|
||||||
).await {
|
).await {
|
||||||
tracing::error!("handle_tick error: {}", e);
|
tracing::error!("handle_tick error: {}", e);
|
||||||
}
|
}
|
||||||
|
|
@ -569,7 +559,6 @@ async fn handle_udp_packet(
|
||||||
peer_last_seen: &mut HashMap<IpAddr, Instant>,
|
peer_last_seen: &mut HashMap<IpAddr, Instant>,
|
||||||
peer_available: &mut HashMap<IpAddr, bool>,
|
peer_available: &mut HashMap<IpAddr, bool>,
|
||||||
last_empty_app_log: &mut Instant,
|
last_empty_app_log: &mut Instant,
|
||||||
session_backpressure: &SessionBackpressure,
|
|
||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
let size = packet.len();
|
let size = packet.len();
|
||||||
match dispatcher.on_datagram(peer, packet.clone()) {
|
match dispatcher.on_datagram(peer, packet.clone()) {
|
||||||
|
|
@ -632,7 +621,6 @@ async fn handle_udp_packet(
|
||||||
connect_tx.clone(),
|
connect_tx.clone(),
|
||||||
router.clone(),
|
router.clone(),
|
||||||
tcp_map,
|
tcp_map,
|
||||||
session_backpressure,
|
|
||||||
).await?;
|
).await?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -651,7 +639,6 @@ async fn handle_tick(
|
||||||
ui_event_tx: &mpsc::UnboundedSender<UiEvent>,
|
ui_event_tx: &mpsc::UnboundedSender<UiEvent>,
|
||||||
peer_last_seen: &mut HashMap<IpAddr, Instant>,
|
peer_last_seen: &mut HashMap<IpAddr, Instant>,
|
||||||
peer_available: &mut HashMap<IpAddr, bool>,
|
peer_available: &mut HashMap<IpAddr, bool>,
|
||||||
session_backpressure: &SessionBackpressure,
|
|
||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
let now = Instant::now();
|
let now = Instant::now();
|
||||||
let peer_timeout = Duration::from_secs(45);
|
let peer_timeout = Duration::from_secs(45);
|
||||||
|
|
@ -662,22 +649,6 @@ async fn handle_tick(
|
||||||
let _ = ui_event_tx.send(UiEvent::Log(format!("Client {peer_ip} disconnected (timeout)")));
|
let _ = ui_event_tx.send(UiEvent::Log(format!("Client {peer_ip} disconnected (timeout)")));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// Publish each active session's current download-direction headroom so
|
|
||||||
// relay reader tasks (running on other tasks, no access to `dispatcher`)
|
|
||||||
// can throttle without touching a lock on every read. New sessions get an
|
|
||||||
// entry created here on their first tick after the handshake; entries for
|
|
||||||
// sessions that no longer exist are pruned below alongside dropped_sessions.
|
|
||||||
{
|
|
||||||
let snapshot = dispatcher.snapshot_backpressure();
|
|
||||||
let mut map = session_backpressure.write().unwrap_or_else(|e| e.into_inner());
|
|
||||||
for (sid, available) in snapshot {
|
|
||||||
match map.get(&sid) {
|
|
||||||
Some(slot) => slot.store(available, std::sync::atomic::Ordering::Relaxed),
|
|
||||||
None => { map.insert(sid, Arc::new(AtomicI64::new(available))); }
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
let (frames, dropped_sessions) = dispatcher.on_tick();
|
let (frames, dropped_sessions) = dispatcher.on_tick();
|
||||||
for (frame, peer_addr) in frames {
|
for (frame, peer_addr) in frames {
|
||||||
let mut sent_tcp = false;
|
let mut sent_tcp = false;
|
||||||
|
|
@ -692,12 +663,6 @@ async fn handle_tick(
|
||||||
let _ = socket.send_to(&frame, peer_addr).await?;
|
let _ = socket.send_to(&frame, peer_addr).await?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if !dropped_sessions.is_empty() {
|
|
||||||
let mut map = session_backpressure.write().unwrap_or_else(|e| e.into_inner());
|
|
||||||
for sid in &dropped_sessions {
|
|
||||||
map.remove(sid);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
for sid in dropped_sessions {
|
for sid in dropped_sessions {
|
||||||
let _ = ui_event_tx.send(UiEvent::Log(format!("Session {sid} expired, releasing resources")));
|
let _ = ui_event_tx.send(UiEvent::Log(format!("Session {sid} expired, releasing resources")));
|
||||||
let mut streams_to_cancel = Vec::new();
|
let mut streams_to_cancel = Vec::new();
|
||||||
|
|
|
||||||
|
|
@ -51,63 +51,19 @@ pub async fn connect_target(
|
||||||
return match outbound.protocol.as_str() {
|
return match outbound.protocol.as_str() {
|
||||||
"socks5" => connect_via_socks5(&proxy_addr, target).await,
|
"socks5" => connect_via_socks5(&proxy_addr, target).await,
|
||||||
"http" => connect_via_http(&proxy_addr, target).await,
|
"http" => connect_via_http(&proxy_addr, target).await,
|
||||||
_ => connect_direct(target, connect_timeout).await,
|
_ => tokio::time::timeout(connect_timeout, TcpStream::connect(target))
|
||||||
|
.await
|
||||||
|
.map_err(|_| anyhow::anyhow!("connect timeout ({}s): {}", connect_timeout.as_secs(), target))?
|
||||||
|
.map_err(Into::into),
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
connect_direct(target, connect_timeout).await
|
tokio::time::timeout(connect_timeout, TcpStream::connect(target))
|
||||||
}
|
.await
|
||||||
|
.map_err(|_| anyhow::anyhow!("connect timeout ({}s): {}", connect_timeout.as_secs(), target))?
|
||||||
/// Per-candidate-address connect attempt, tried in turn (see `connect_direct`
|
.map_err(Into::into)
|
||||||
/// below). Short enough that a single dead-end address can't eat the whole
|
|
||||||
/// outer `connect_timeout` budget.
|
|
||||||
const PER_ADDR_CONNECT_TIMEOUT: Duration = Duration::from_secs(3);
|
|
||||||
|
|
||||||
/// Resolve `target` ("host:port") and connect to it, trying candidate
|
|
||||||
/// addresses in turn rather than handing the raw string straight to
|
|
||||||
/// `TcpStream::connect` (which resolves and tries addresses internally but
|
|
||||||
/// shares ONE timeout across the whole attempt).
|
|
||||||
///
|
|
||||||
/// IPv4 candidates are tried first. Some VPS hosts (observed on a
|
|
||||||
/// DigitalOcean droplet) assign the machine an IPv6 address that the OS
|
|
||||||
/// prefers by RFC 6724 ordering but that has no actually-working outbound
|
|
||||||
/// route - the connect attempt doesn't get refused, it just hangs. With a
|
|
||||||
/// single shared timeout across all candidates, that one dead IPv6 address
|
|
||||||
/// eats the entire budget and the working IPv4 candidate is never even
|
|
||||||
/// attempted: every dual-stack destination (i.e. most popular sites) never
|
|
||||||
/// loads, while IPv4-only destinations work fine - exactly the "traffic
|
|
||||||
/// counter moves but sites don't open" symptom this fixes.
|
|
||||||
async fn connect_direct(target: &str, connect_timeout: Duration) -> Result<TcpStream> {
|
|
||||||
tokio::time::timeout(connect_timeout, async {
|
|
||||||
let mut addrs: Vec<std::net::SocketAddr> = tokio::net::lookup_host(target)
|
|
||||||
.await
|
|
||||||
.map_err(|e| anyhow::anyhow!("dns resolution failed for {}: {}", target, e))?
|
|
||||||
.collect();
|
|
||||||
if addrs.is_empty() {
|
|
||||||
return Err(anyhow::anyhow!("no addresses resolved for {}", target));
|
|
||||||
}
|
|
||||||
prefer_ipv4_first(&mut addrs);
|
|
||||||
|
|
||||||
let mut last_err = None;
|
|
||||||
for addr in addrs {
|
|
||||||
match tokio::time::timeout(PER_ADDR_CONNECT_TIMEOUT, TcpStream::connect(addr)).await {
|
|
||||||
Ok(Ok(stream)) => return Ok(stream),
|
|
||||||
Ok(Err(e)) => last_err = Some(anyhow::anyhow!("{}: {}", addr, e)),
|
|
||||||
Err(_) => last_err = Some(anyhow::anyhow!("{}: connect timeout ({}s)", addr, PER_ADDR_CONNECT_TIMEOUT.as_secs())),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Err(last_err.unwrap_or_else(|| anyhow::anyhow!("all candidates failed for {}", target)))
|
|
||||||
})
|
|
||||||
.await
|
|
||||||
.map_err(|_| anyhow::anyhow!("connect timeout ({}s): {}", connect_timeout.as_secs(), target))?
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Stable-sort so IPv4 candidates come before IPv6 ones, without otherwise
|
|
||||||
/// disturbing the resolver's original ordering within each family.
|
|
||||||
fn prefer_ipv4_first(addrs: &mut [std::net::SocketAddr]) {
|
|
||||||
addrs.sort_by_key(|a| a.is_ipv6());
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// ── Rule matching ────────────────────────────────────────────────────────────
|
// ── Rule matching ────────────────────────────────────────────────────────────
|
||||||
|
|
@ -584,49 +540,4 @@ mod tests {
|
||||||
fn test_match_domain_rule_empty() {
|
fn test_match_domain_rule_empty() {
|
||||||
assert!(!match_domain_rule("example.com", &[]));
|
assert!(!match_domain_rule("example.com", &[]));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_prefer_ipv4_first_reorders_mixed_list() {
|
|
||||||
let v6: std::net::SocketAddr = "[2001:db8::1]:443".parse().unwrap();
|
|
||||||
let v4: std::net::SocketAddr = "192.0.2.1:443".parse().unwrap();
|
|
||||||
let mut addrs = vec![v6, v4];
|
|
||||||
prefer_ipv4_first(&mut addrs);
|
|
||||||
assert_eq!(addrs, vec![v4, v6], "IPv4 candidate must sort before IPv6");
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_prefer_ipv4_first_preserves_order_within_family() {
|
|
||||||
// Two IPv4 addresses: relative order should be untouched (stable sort).
|
|
||||||
let a: std::net::SocketAddr = "192.0.2.1:443".parse().unwrap();
|
|
||||||
let b: std::net::SocketAddr = "192.0.2.2:443".parse().unwrap();
|
|
||||||
let mut addrs = vec![a, b];
|
|
||||||
prefer_ipv4_first(&mut addrs);
|
|
||||||
assert_eq!(addrs, vec![a, b]);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tokio::test]
|
|
||||||
async fn test_connect_direct_succeeds_against_live_listener() {
|
|
||||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
|
||||||
let addr = listener.local_addr().unwrap();
|
|
||||||
tokio::spawn(async move {
|
|
||||||
let _ = listener.accept().await;
|
|
||||||
});
|
|
||||||
|
|
||||||
let result = connect_direct(&addr.to_string(), Duration::from_secs(2)).await;
|
|
||||||
assert!(result.is_ok(), "expected connect_direct to reach a live local listener: {:?}", result.err());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tokio::test]
|
|
||||||
async fn test_connect_direct_fails_fast_on_refused_port() {
|
|
||||||
// Bind and immediately drop to get a port nothing is listening on,
|
|
||||||
// so the OS sends RST and the attempt fails well under the timeout.
|
|
||||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
|
||||||
let addr = listener.local_addr().unwrap();
|
|
||||||
drop(listener);
|
|
||||||
|
|
||||||
let start = std::time::Instant::now();
|
|
||||||
let result = connect_direct(&addr.to_string(), Duration::from_secs(5)).await;
|
|
||||||
assert!(result.is_err(), "connecting to a closed port should fail");
|
|
||||||
assert!(start.elapsed() < Duration::from_secs(4), "a refused connection must not wait out the full timeout");
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,8 +1,6 @@
|
||||||
use anyhow::Result;
|
use anyhow::Result;
|
||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
use portable_atomic::AtomicI64;
|
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::sync::Arc;
|
|
||||||
|
|
||||||
use ostp_core::relay::RelayMessage;
|
use ostp_core::relay::RelayMessage;
|
||||||
use tokio::io::AsyncReadExt;
|
use tokio::io::AsyncReadExt;
|
||||||
|
|
@ -10,19 +8,7 @@ use tokio::net::UdpSocket;
|
||||||
use tokio::sync::mpsc;
|
use tokio::sync::mpsc;
|
||||||
|
|
||||||
use crate::dispatcher::Dispatcher;
|
use crate::dispatcher::Dispatcher;
|
||||||
use crate::{RemoteState, SessionBackpressure, UiEvent};
|
use crate::{RemoteState, UiEvent};
|
||||||
|
|
||||||
/// How long a target-connection reader task waits before rechecking the
|
|
||||||
/// client session's congestion headroom while throttled. Short enough that
|
|
||||||
/// a freed-up window (checked every server tick, 10ms) is noticed promptly;
|
|
||||||
/// long enough not to spin.
|
|
||||||
const BACKPRESSURE_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_millis(5);
|
|
||||||
/// Upper bound on total time a single read is throttled before proceeding
|
|
||||||
/// anyway. Congestion state is a hint, not a hard guarantee - if the
|
|
||||||
/// session's headroom never frees up (e.g. a stuck/buggy state), a stream
|
|
||||||
/// must not be stalled forever; better to occasionally overshoot the window
|
|
||||||
/// than deadlock a connection.
|
|
||||||
const BACKPRESSURE_MAX_WAIT: std::time::Duration = std::time::Duration::from_secs(2);
|
|
||||||
|
|
||||||
fn clean_ipv6_mapped_v4(addr: std::net::SocketAddr) -> std::net::SocketAddr {
|
fn clean_ipv6_mapped_v4(addr: std::net::SocketAddr) -> std::net::SocketAddr {
|
||||||
match addr {
|
match addr {
|
||||||
|
|
@ -52,7 +38,6 @@ pub async fn handle_relay_message(
|
||||||
connect_tx: mpsc::UnboundedSender<(u32, u16, String, Result<(tokio::net::tcp::OwnedWriteHalf, mpsc::Sender<()>), String>)>,
|
connect_tx: mpsc::UnboundedSender<(u32, u16, String, Result<(tokio::net::tcp::OwnedWriteHalf, mpsc::Sender<()>), String>)>,
|
||||||
router: std::sync::Arc<crate::router::Router>,
|
router: std::sync::Arc<crate::router::Router>,
|
||||||
tcp_map: &std::sync::Arc<tokio::sync::RwLock<HashMap<std::net::SocketAddr, tokio::sync::mpsc::Sender<Bytes>>>>,
|
tcp_map: &std::sync::Arc<tokio::sync::RwLock<HashMap<std::net::SocketAddr, tokio::sync::mpsc::Sender<Bytes>>>>,
|
||||||
session_backpressure: &SessionBackpressure,
|
|
||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
match RelayMessage::decode(&payload)? {
|
match RelayMessage::decode(&payload)? {
|
||||||
RelayMessage::Connect(target) => {
|
RelayMessage::Connect(target) => {
|
||||||
|
|
@ -68,41 +53,15 @@ pub async fn handle_relay_message(
|
||||||
let connect_tx_clone = connect_tx.clone();
|
let connect_tx_clone = connect_tx.clone();
|
||||||
let stream_tx_clone = stream_tx.clone();
|
let stream_tx_clone = stream_tx.clone();
|
||||||
let router_clone = router.clone();
|
let router_clone = router.clone();
|
||||||
let backpressure_clone = session_backpressure.clone();
|
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
let stream_res = router_clone.route_tcp(&target_clone).await;
|
let stream_res = router_clone.route_tcp(&target_clone).await;
|
||||||
match stream_res {
|
match stream_res {
|
||||||
Ok(stream) => {
|
Ok(stream) => {
|
||||||
let (mut reader, writer) = stream.into_split();
|
let (mut reader, writer) = stream.into_split();
|
||||||
let (cancel_tx, mut cancel_rx) = mpsc::channel::<()>(1);
|
let (cancel_tx, mut cancel_rx) = mpsc::channel::<()>(1);
|
||||||
// Get-or-create this session's headroom handle. A brand
|
|
||||||
// new session may not have its first tick's snapshot
|
|
||||||
// yet (up to 10ms), so default it open (matches a fresh
|
|
||||||
// congestion window) rather than stalling the very
|
|
||||||
// first read while nothing has been published.
|
|
||||||
let headroom: Arc<AtomicI64> = {
|
|
||||||
let mut map = backpressure_clone.write().unwrap_or_else(|e| e.into_inner());
|
|
||||||
map.entry(session_id).or_insert_with(|| Arc::new(AtomicI64::new(32))).clone()
|
|
||||||
};
|
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
let mut buf = [0_u8; 4096];
|
let mut buf = [0_u8; 4096];
|
||||||
loop {
|
loop {
|
||||||
// Throttle to the client-facing OSTP session's
|
|
||||||
// real congestion window instead of reading from
|
|
||||||
// the target as fast as it'll send. Without this,
|
|
||||||
// a fast target blasts a lossy/jittery client
|
|
||||||
// path far beyond what it can sustain, which
|
|
||||||
// self-inflicts a loss burst, wrecks the RTT/RTO
|
|
||||||
// estimate, and can stall the session hard
|
|
||||||
// enough to trip the client's keepalive
|
|
||||||
// reconnect. See Dispatcher::snapshot_backpressure.
|
|
||||||
let mut waited = std::time::Duration::ZERO;
|
|
||||||
while headroom.load(std::sync::atomic::Ordering::Relaxed) <= 0
|
|
||||||
&& waited < BACKPRESSURE_MAX_WAIT
|
|
||||||
{
|
|
||||||
tokio::time::sleep(BACKPRESSURE_POLL_INTERVAL).await;
|
|
||||||
waited += BACKPRESSURE_POLL_INTERVAL;
|
|
||||||
}
|
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
_ = cancel_rx.recv() => break,
|
_ = cancel_rx.recv() => break,
|
||||||
read_res = reader.read(&mut buf) => {
|
read_res = reader.read(&mut buf) => {
|
||||||
|
|
|
||||||
|
|
@ -138,5 +138,5 @@ Write-Host "No configuration found. Launching setup wizard..."
|
||||||
Write-Host ""
|
Write-Host ""
|
||||||
|
|
||||||
Push-Location $InstallDir
|
Push-Location $InstallDir
|
||||||
& .\ostp.exe setup
|
& .\ostp.exe --setup
|
||||||
Pop-Location
|
Pop-Location
|
||||||
|
|
|
||||||
|
|
@ -238,4 +238,4 @@ echo "No configuration found. Launching setup wizard..."
|
||||||
echo ""
|
echo ""
|
||||||
|
|
||||||
cd "$INSTALL_DIR"
|
cd "$INSTALL_DIR"
|
||||||
exec ./ostp setup --config "$CONFIG_FILE"
|
exec ./ostp --setup --config "$CONFIG_FILE"
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue