mirror of https://github.com/ospab/ostp.git
Compare commits
13 Commits
cd12b01bc3
...
cdfd2babc0
| Author | SHA1 | Date |
|---|---|---|
|
|
cdfd2babc0 | |
|
|
2092e22a7c | |
|
|
5278f58903 | |
|
|
340819745a | |
|
|
e31c4b2268 | |
|
|
e46c863ef0 | |
|
|
cddd623ad0 | |
|
|
9a891310f9 | |
|
|
d9686c9344 | |
|
|
dbf923fb16 | |
|
|
51b947e6ff | |
|
|
f01ed4ec25 | |
|
|
c2a1a53b4d |
|
|
@ -284,7 +284,15 @@ jobs:
|
|||
|
||||
- name: Install cross (if not cached)
|
||||
if: ${{ matrix.use_cross && steps.cross-cache.outputs.cache-hit != 'true' }}
|
||||
run: cargo install cross --git https://github.com/cross-rs/cross.git --locked
|
||||
# cross-rs's own source (not ours, not a dependency of ours) uses a
|
||||
# macro-at-end-of-block pattern that trips rustc's
|
||||
# semicolon_in_expressions_from_macros lint on current toolchains -
|
||||
# harmless in cross's actual behavior, but `cargo install` compiles
|
||||
# the installed package as the "local" crate, so dependency lint
|
||||
# capping doesn't shield it. --cap-lints=warn is the standard escape
|
||||
# hatch for building a third-party tool against a newer compiler than
|
||||
# its own lint config assumed; it doesn't touch our own build.
|
||||
run: RUSTFLAGS="--cap-lints=warn" cargo install cross --git https://github.com/cross-rs/cross.git --locked
|
||||
|
||||
- name: Build (cross)
|
||||
if: ${{ matrix.use_cross }}
|
||||
|
|
|
|||
|
|
@ -2,5 +2,5 @@
|
|||
"target_version": "0.4.2",
|
||||
"branch": "beta",
|
||||
"alpha_iteration": 0,
|
||||
"beta_iteration": 2
|
||||
"beta_iteration": 4
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1400,6 +1400,7 @@ dependencies = [
|
|||
"rlimit",
|
||||
"serde",
|
||||
"serde_json",
|
||||
"sha2",
|
||||
"tokio",
|
||||
"tracing",
|
||||
"tracing-subscriber",
|
||||
|
|
@ -1496,6 +1497,7 @@ dependencies = [
|
|||
"sha2",
|
||||
"simple-dns",
|
||||
"socket2",
|
||||
"subtle",
|
||||
"tokio",
|
||||
"tower-http",
|
||||
"tracing",
|
||||
|
|
|
|||
|
|
@ -231,7 +231,7 @@ impl Bridge {
|
|||
self.handle_inbound_udp(udp_msg, &mut sessions_opt, &mut udp_rx_opt, &mut proxy_guard, &mut stream_map, &tx, &proxy_tx).await;
|
||||
}
|
||||
cmd = bridge_rx.recv() => {
|
||||
if !self.handle_bridge_cmd(cmd, &mut sessions_opt, &mut udp_rx_opt, &mut proxy_guard, &mut stream_map, &tx, &proxy_tx).await {
|
||||
if !self.handle_bridge_cmd(cmd, &mut bridge_rx, &mut sessions_opt, &mut udp_rx_opt, &mut proxy_guard, &mut stream_map, &tx, &proxy_tx).await {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
|
@ -374,6 +374,7 @@ impl Bridge {
|
|||
async fn handle_bridge_cmd(
|
||||
&mut self,
|
||||
cmd: Option<BridgeCommand>,
|
||||
bridge_rx: &mut mpsc::Receiver<BridgeCommand>,
|
||||
sessions_opt: &mut Option<Vec<SessionState>>,
|
||||
udp_rx_opt: &mut Option<mpsc::Receiver<(usize, Bytes)>>,
|
||||
proxy_guard: &mut Option<crate::sysproxy::SystemProxyGuard>,
|
||||
|
|
@ -465,6 +466,32 @@ impl Bridge {
|
|||
tx.send(UiEvent::Log(format!("Obfuscation profile switched to {:?}", self.profile))).await.ok();
|
||||
}
|
||||
Some(BridgeCommand::NetworkChanged) => {
|
||||
// A real network handoff (Wi-Fi <-> cellular) commonly fires
|
||||
// onLost + onAvailable within milliseconds of each other on
|
||||
// Android, queuing several NetworkChanged commands back to
|
||||
// back. Each reconnect below is a full sequential handshake
|
||||
// (up to ~1.2s x 4 attempts x mux_sessions) run synchronously
|
||||
// in this select-loop iteration, so without coalescing, the
|
||||
// first attempt often races the OS's own network switch and
|
||||
// fails on the now-dead interface, then the SECOND queued
|
||||
// NetworkChanged only starts its own full reconnect after
|
||||
// that first one finishes - multiplying a sub-second handoff
|
||||
// into many seconds of extra outage. Drain same-kind repeats
|
||||
// so a burst collapses into one reconnect on the freshest
|
||||
// signal; a different command found while draining is
|
||||
// handled immediately rather than dropped.
|
||||
while let Ok(next) = bridge_rx.try_recv() {
|
||||
if !matches!(next, BridgeCommand::NetworkChanged) {
|
||||
let more = Box::pin(self.handle_bridge_cmd(
|
||||
Some(next), bridge_rx, sessions_opt, udp_rx_opt, proxy_guard, stream_map, tx, proxy_tx,
|
||||
)).await;
|
||||
if !more {
|
||||
return false;
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if self.running {
|
||||
let _ = tx.send(UiEvent::Log("Network changed — starting immediate reconnect".to_string())).await;
|
||||
self.metrics.connection_state.store(1, Ordering::Relaxed);
|
||||
|
|
|
|||
|
|
@ -361,6 +361,10 @@ async fn handle_udp_associate(
|
|||
|
||||
let mut direct_udp_v4: Option<Arc<UdpSocket>> = None;
|
||||
let mut direct_udp_v6: Option<Arc<UdpSocket>> = None;
|
||||
// Held only to keep the direct-UDP readers' cancellation senders alive;
|
||||
// dropping this (on every return path from this function) is what tells
|
||||
// spawn_direct_udp_reader's tasks to stop. See its doc comment.
|
||||
let mut direct_udp_cancel_txs: Vec<tokio::sync::oneshot::Sender<()>> = Vec::new();
|
||||
|
||||
let mut tcp_buf = [0u8; 1];
|
||||
loop {
|
||||
|
|
@ -432,7 +436,9 @@ async fn handle_udp_associate(
|
|||
match create_udp_socket_bypassing_tun(true, matcher.physical_if_index, &matcher.physical_if_name).await {
|
||||
Ok(s) => {
|
||||
let s_arc = Arc::new(s);
|
||||
spawn_direct_udp_reader(s_arc.clone(), sock_tx.clone(), client_udp_addr.clone(), debug);
|
||||
let (cancel_tx, cancel_rx) = tokio::sync::oneshot::channel();
|
||||
spawn_direct_udp_reader(s_arc.clone(), sock_tx.clone(), client_udp_addr.clone(), debug, cancel_rx);
|
||||
direct_udp_cancel_txs.push(cancel_tx);
|
||||
direct_udp_v6 = Some(s_arc);
|
||||
}
|
||||
Err(e) => {
|
||||
|
|
@ -446,7 +452,9 @@ async fn handle_udp_associate(
|
|||
match create_udp_socket_bypassing_tun(false, matcher.physical_if_index, &matcher.physical_if_name).await {
|
||||
Ok(s) => {
|
||||
let s_arc = Arc::new(s);
|
||||
spawn_direct_udp_reader(s_arc.clone(), sock_tx.clone(), client_udp_addr.clone(), debug);
|
||||
let (cancel_tx, cancel_rx) = tokio::sync::oneshot::channel();
|
||||
spawn_direct_udp_reader(s_arc.clone(), sock_tx.clone(), client_udp_addr.clone(), debug, cancel_rx);
|
||||
direct_udp_cancel_txs.push(cancel_tx);
|
||||
direct_udp_v4 = Some(s_arc);
|
||||
}
|
||||
Err(e) => {
|
||||
|
|
@ -520,11 +528,24 @@ fn spawn_direct_udp_reader(
|
|||
sock_tx: Arc<UdpSocket>,
|
||||
client_udp_addr: Arc<std::sync::Mutex<Option<std::net::SocketAddr>>>,
|
||||
_debug: bool,
|
||||
mut cancel_rx: tokio::sync::oneshot::Receiver<()>,
|
||||
) {
|
||||
tokio::spawn(async move {
|
||||
let mut buf = vec![0u8; 65536];
|
||||
loop {
|
||||
match direct_socket.recv_from(&mut buf).await {
|
||||
let recv_result = tokio::select! {
|
||||
// Fires as soon as the sender half (held by handle_udp_associate
|
||||
// for exactly this reason) is dropped - which happens the
|
||||
// instant that function returns, on every exit path, with no
|
||||
// explicit signaling needed. Without this, a UDP-associate
|
||||
// session that ever bypassed traffic direct (excluded IP/
|
||||
// domain) leaked this socket + task for the rest of the
|
||||
// process's life once the session ended: nothing else ever
|
||||
// stopped this loop.
|
||||
_ = &mut cancel_rx => break,
|
||||
res = direct_socket.recv_from(&mut buf) => res,
|
||||
};
|
||||
match recv_result {
|
||||
Ok((len, target_addr)) => {
|
||||
let client_addr = {
|
||||
let guard = client_udp_addr.lock().unwrap();
|
||||
|
|
|
|||
|
|
@ -138,27 +138,34 @@ async fn start_udp_bypass_session(
|
|||
let _ = crate::tunnel::proxy::bind_socket_to_interface(&socket, name);
|
||||
}
|
||||
|
||||
let socket = Arc::new(socket);
|
||||
let socket_rx = socket.clone();
|
||||
|
||||
// Spawn a task to read from physical socket and send back to smoltcp
|
||||
let tx_clone = smoltcp_tx.clone();
|
||||
tokio::spawn(async move {
|
||||
// A single select! loop over both directions, rather than spawning a
|
||||
// separate task for the read side, so the whole session - physical
|
||||
// socket included - is torn down the moment this function returns
|
||||
// (e.g. when session_rx closes). The previous spawned-task version left
|
||||
// that task (and its Arc<UdpSocket> clone, keeping the OS socket fd
|
||||
// alive) running forever after this function returned: nothing ever
|
||||
// cancelled it, so every bypassed UDP flow (any excluded app/IP in TUN
|
||||
// mode) leaked one socket + one task for the lifetime of the process.
|
||||
use futures::SinkExt;
|
||||
let mut buf = [0u8; 65536];
|
||||
loop {
|
||||
match socket_rx.recv_from(&mut buf).await {
|
||||
tokio::select! {
|
||||
outbound = session_rx.recv() => {
|
||||
match outbound {
|
||||
Some((payload, dst)) => { socket.send_to(&payload, dst).await?; }
|
||||
None => break,
|
||||
}
|
||||
}
|
||||
inbound = socket.recv_from(&mut buf) => {
|
||||
match inbound {
|
||||
Ok((n, peer)) => {
|
||||
let mut lock = tx_clone.lock().await;
|
||||
let mut lock = smoltcp_tx.lock().await;
|
||||
let _ = lock.send((buf[..n].to_vec(), peer, client_src)).await;
|
||||
}
|
||||
Err(_) => break,
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
while let Some((payload, dst)) = session_rx.recv().await {
|
||||
socket.send_to(&payload, dst).await?;
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
|
|
|
|||
|
|
@ -43,6 +43,11 @@ pub struct CongestionController {
|
|||
mtu: u64,
|
||||
/// Min RTT expiry: re-probe after 10 seconds
|
||||
min_rtt_stamp: Instant,
|
||||
/// Loss events counted toward SLOW_START_LOSS_TOLERANCE within the
|
||||
/// current SLOW_START_LOSS_WINDOW (see on_loss's SlowStart arm).
|
||||
slow_start_losses: u32,
|
||||
/// Start of the current loss-tolerance window.
|
||||
slow_start_loss_window_start: Instant,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
|
|
@ -67,6 +72,24 @@ const RTO_MAX: Duration = Duration::from_secs(16);
|
|||
/// Will be replaced by first real measurement within milliseconds.
|
||||
const INITIAL_RTT: Duration = Duration::from_millis(30);
|
||||
|
||||
/// Isolated packet loss during slow start (a single dropped frame from
|
||||
/// wireless noise, a brief LTE handover blip, etc.) is normal on real
|
||||
/// mobile/Wi-Fi links and does NOT mean the link is congested. The previous
|
||||
/// behavior exited slow start and halved cwnd on the very FIRST loss, which
|
||||
/// on any link with a non-zero background loss rate permanently downgrades
|
||||
/// the session from exponential growth to linear (+1 MTU/RTT) ProbeBandwidth
|
||||
/// growth within the first few RTTs - turning what should be a sub-second
|
||||
/// ramp-up into tens of seconds to minutes before throughput opens up
|
||||
/// (observed as: a trickle of KB/s, then a sudden jump once cwnd finally
|
||||
/// claws back up). Only treat loss as a real congestion signal - and pay
|
||||
/// the full slow-start-exit + halving cost - once this many losses land
|
||||
/// within SLOW_START_LOSS_WINDOW.
|
||||
const SLOW_START_LOSS_TOLERANCE: u32 = 3;
|
||||
/// Window within which SLOW_START_LOSS_TOLERANCE losses must land to count
|
||||
/// as sustained (rather than isolated) loss. Roughly a few RTTs on a
|
||||
/// well-connected link, generous on a slow one.
|
||||
const SLOW_START_LOSS_WINDOW: Duration = Duration::from_millis(500);
|
||||
|
||||
impl CongestionController {
|
||||
pub fn new(mtu: u64) -> Self {
|
||||
let now = Instant::now();
|
||||
|
|
@ -88,6 +111,8 @@ impl CongestionController {
|
|||
pacing_rate: initial_pacing,
|
||||
mtu,
|
||||
min_rtt_stamp: now,
|
||||
slow_start_losses: 0,
|
||||
slow_start_loss_window_start: now,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -197,11 +222,28 @@ impl CongestionController {
|
|||
|
||||
match self.phase {
|
||||
Phase::SlowStart => {
|
||||
// Exit slow start, set ssthresh to half of cwnd
|
||||
let now = Instant::now();
|
||||
if now.duration_since(self.slow_start_loss_window_start) > SLOW_START_LOSS_WINDOW {
|
||||
// Previous window's losses have aged out - this loss starts a fresh count.
|
||||
self.slow_start_losses = 0;
|
||||
self.slow_start_loss_window_start = now;
|
||||
}
|
||||
self.slow_start_losses += 1;
|
||||
|
||||
if self.slow_start_losses >= SLOW_START_LOSS_TOLERANCE {
|
||||
// Sustained loss within the window: treat as real congestion.
|
||||
// Exit slow start, set ssthresh to half of cwnd.
|
||||
self.ssthresh = self.cwnd / 2;
|
||||
self.cwnd = self.ssthresh.max(MIN_CWND_PACKETS * self.mtu);
|
||||
self.phase = Phase::ProbeBandwidth;
|
||||
tracing::debug!(cwnd = self.cwnd, ssthresh = self.ssthresh, "congestion: loss during slow start");
|
||||
tracing::debug!(cwnd = self.cwnd, ssthresh = self.ssthresh, "congestion: sustained loss during slow start, exiting");
|
||||
} else {
|
||||
// Isolated loss: likely non-congestive noise. Take a mild,
|
||||
// temporary haircut but keep exponential growth going -
|
||||
// don't throw away slow start over a single dropped frame.
|
||||
self.cwnd = (self.cwnd * 8 / 10).max(MIN_CWND_PACKETS * self.mtu);
|
||||
tracing::debug!(cwnd = self.cwnd, count = self.slow_start_losses, "congestion: isolated loss during slow start, staying in slow start");
|
||||
}
|
||||
}
|
||||
Phase::ProbeBandwidth => {
|
||||
// Multiplicative decrease: cwnd *= 0.7 (BBR-style, less aggressive than Cubic's 0.5)
|
||||
|
|
@ -290,6 +332,50 @@ mod tests {
|
|||
assert!(cc.cwnd() < initial);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_isolated_slow_start_loss_does_not_exit_slow_start() {
|
||||
// A single dropped packet (wireless noise, a brief handover blip) is
|
||||
// normal on real links and must not permanently downgrade the
|
||||
// session from exponential to linear growth.
|
||||
let mut cc = CongestionController::new(1200);
|
||||
cc.on_loss(1200);
|
||||
assert_eq!(cc.phase, Phase::SlowStart, "one isolated loss must not exit slow start");
|
||||
|
||||
// It should still shrink the window somewhat (not ignored entirely),
|
||||
// just far less punishing than the sustained-congestion case.
|
||||
let after_one = cc.cwnd();
|
||||
assert!(after_one < INITIAL_CWND_PACKETS * 1200);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_sustained_slow_start_loss_exits_slow_start() {
|
||||
// Losses landing close together (within SLOW_START_LOSS_WINDOW) are
|
||||
// a real congestion signal and must still trigger the harsher
|
||||
// exit-slow-start + halve response.
|
||||
let mut cc = CongestionController::new(1200);
|
||||
for _ in 0..SLOW_START_LOSS_TOLERANCE {
|
||||
cc.on_loss(1200);
|
||||
}
|
||||
assert_eq!(cc.phase, Phase::ProbeBandwidth, "sustained loss must exit slow start");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_slow_start_loss_window_resets_after_expiry() {
|
||||
// Two losses far enough apart (window expired between them) must
|
||||
// each be treated as isolated, not accumulated toward the sustained-
|
||||
// loss threshold.
|
||||
let mut cc = CongestionController::new(1200);
|
||||
cc.on_loss(1200);
|
||||
assert_eq!(cc.phase, Phase::SlowStart);
|
||||
|
||||
// Simulate the window having expired by resetting its start
|
||||
// directly (std::thread::sleep in a unit test would be flaky/slow).
|
||||
cc.slow_start_loss_window_start = Instant::now() - SLOW_START_LOSS_WINDOW - Duration::from_millis(1);
|
||||
cc.on_loss(1200);
|
||||
assert_eq!(cc.phase, Phase::SlowStart, "a loss after the window expired must restart the count, not accumulate");
|
||||
assert_eq!(cc.slow_start_losses, 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_can_send_limits() {
|
||||
let mut cc = CongestionController::new(1200);
|
||||
|
|
|
|||
|
|
@ -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
|
||||
# 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.
|
||||
version: 0.4.2+21
|
||||
version: 0.4.2+23
|
||||
|
||||
environment:
|
||||
sdk: ^3.11.4
|
||||
|
|
|
|||
|
|
@ -31,3 +31,4 @@ hex = "0.4.3"
|
|||
chacha20poly1305.workspace = true
|
||||
x25519-dalek = { version = "2.0.1", features = ["static_secrets"] }
|
||||
chrono = "0.4.44"
|
||||
subtle = "2.6"
|
||||
|
|
|
|||
|
|
@ -318,6 +318,18 @@ pub async fn start_api_server(
|
|||
|
||||
// ── Middleware: token check ──────────────────────────────────────────────────
|
||||
|
||||
/// Constant-time string equality for secrets (tokens, password hashes).
|
||||
/// Plain `==` short-circuits on the first differing byte, which leaks how
|
||||
/// many leading bytes an attacker's guess got right through response
|
||||
/// timing - a classic remote timing side-channel against exactly the kind
|
||||
/// of long-lived bearer/session secrets compared here. `subtle` is already
|
||||
/// pulled in transitively (chacha20poly1305 etc.); pinning it as a direct
|
||||
/// dependency here makes that guarantee explicit for this call site.
|
||||
fn secure_eq(a: &str, b: &str) -> bool {
|
||||
use subtle::ConstantTimeEq;
|
||||
a.as_bytes().ct_eq(b.as_bytes()).into()
|
||||
}
|
||||
|
||||
fn check_token(state: &ApiState, headers: &axum::http::HeaderMap) -> bool {
|
||||
// Both session token (for web UI) and static API token (for relays) are checked
|
||||
let mut allowed = false;
|
||||
|
|
@ -332,19 +344,19 @@ fn check_token(state: &ApiState, headers: &axum::http::HeaderMap) -> bool {
|
|||
if let Some(token) = val.strip_prefix("Bearer ") {
|
||||
let current_session = state.session_token.read().unwrap_or_else(|e| e.into_inner()).clone();
|
||||
if let Some(session) = current_session {
|
||||
if token == session {
|
||||
if secure_eq(token, &session) {
|
||||
allowed = true;
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(ref api_tok) = state.api_token {
|
||||
if token == api_tok {
|
||||
if secure_eq(token, api_tok) {
|
||||
allowed = true;
|
||||
}
|
||||
}
|
||||
} else {
|
||||
if let Some(ref api_tok) = state.api_token {
|
||||
if val == api_tok {
|
||||
if secure_eq(val, api_tok) {
|
||||
allowed = true;
|
||||
}
|
||||
}
|
||||
|
|
@ -371,7 +383,7 @@ async fn handle_login(
|
|||
let hash = sha2::Sha256::digest(password.as_bytes());
|
||||
let hash_hex = format!("{:x}", hash);
|
||||
|
||||
if hash_hex == state.password_hash {
|
||||
if secure_eq(&hash_hex, &state.password_hash) {
|
||||
let token = uuid::Uuid::new_v4().to_string();
|
||||
*state.session_token.write().unwrap_or_else(|e| e.into_inner()) = Some(token.clone());
|
||||
(StatusCode::OK, ApiResponse::success(LoginResponse { token }))
|
||||
|
|
@ -881,15 +893,91 @@ mod tests {
|
|||
let state = make_test_state("");
|
||||
let _router = create_api_router(state);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_secure_eq_matches_and_rejects() {
|
||||
assert!(secure_eq("same-secret", "same-secret"));
|
||||
assert!(!secure_eq("same-secret", "different"));
|
||||
assert!(!secure_eq("short", "much-longer-value"));
|
||||
assert!(secure_eq("", ""));
|
||||
}
|
||||
|
||||
fn headers_with_bearer(token: &str) -> axum::http::HeaderMap {
|
||||
let mut h = axum::http::HeaderMap::new();
|
||||
h.insert("authorization", format!("Bearer {token}").parse().unwrap());
|
||||
h
|
||||
}
|
||||
|
||||
// These pin down check_token's behavior directly: it's the single gate
|
||||
// every mutating/sensitive handler (including the audit-log ones - see
|
||||
// the missing-auth fix) relies on, so its logic must be independently
|
||||
// verified rather than only exercised incidentally through handlers.
|
||||
#[test]
|
||||
fn test_check_token_rejects_missing_header_when_configured() {
|
||||
let state = make_test_state("panel");
|
||||
assert!(!check_token(&state, &axum::http::HeaderMap::new()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_check_token_accepts_matching_api_token_as_bearer() {
|
||||
let state = make_test_state("panel");
|
||||
assert!(check_token(&state, &headers_with_bearer("test-token")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_check_token_accepts_matching_api_token_raw() {
|
||||
let state = make_test_state("panel");
|
||||
let mut h = axum::http::HeaderMap::new();
|
||||
h.insert("authorization", "test-token".parse().unwrap());
|
||||
assert!(check_token(&state, &h));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_check_token_rejects_wrong_token() {
|
||||
let state = make_test_state("panel");
|
||||
assert!(!check_token(&state, &headers_with_bearer("wrong-token")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_check_token_accepts_matching_session_token() {
|
||||
let state = make_test_state("panel");
|
||||
*state.session_token.write().unwrap() = Some("live-session".to_string());
|
||||
assert!(check_token(&state, &headers_with_bearer("live-session")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_check_token_open_when_no_credentials_configured() {
|
||||
let mut state = make_test_state("panel");
|
||||
state.api_token = None;
|
||||
state.username.clear();
|
||||
state.password_hash.clear();
|
||||
// Documented "unsafe but possible" open-panel mode: no credentials
|
||||
// configured at all means every request passes, including with no
|
||||
// Authorization header.
|
||||
assert!(check_token(&state, &axum::http::HeaderMap::new()));
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_get_audit(State(state): State<ApiState>) -> impl IntoResponse {
|
||||
let logs = state.audit_logs.read().unwrap();
|
||||
ApiResponse::success(logs.clone())
|
||||
async fn handle_get_audit(
|
||||
State(state): State<ApiState>,
|
||||
headers: axum::http::HeaderMap,
|
||||
) -> impl IntoResponse {
|
||||
if !check_token(&state, &headers) {
|
||||
return api_unauthorized::<Vec<AuditLogEntry>>();
|
||||
}
|
||||
let logs = state.audit_logs.read().unwrap_or_else(|e| e.into_inner());
|
||||
(StatusCode::OK, ApiResponse::success(logs.clone()))
|
||||
}
|
||||
|
||||
async fn handle_create_audit(State(state): State<ApiState>, Json(req): Json<CreateAuditLogRequest>) -> impl IntoResponse {
|
||||
let mut logs = state.audit_logs.write().unwrap();
|
||||
async fn handle_create_audit(
|
||||
State(state): State<ApiState>,
|
||||
headers: axum::http::HeaderMap,
|
||||
Json(req): Json<CreateAuditLogRequest>,
|
||||
) -> impl IntoResponse {
|
||||
if !check_token(&state, &headers) {
|
||||
return api_unauthorized::<bool>();
|
||||
}
|
||||
let mut logs = state.audit_logs.write().unwrap_or_else(|e| e.into_inner());
|
||||
let id = format!("{:x}", rand::random::<u64>());
|
||||
let now = chrono::Local::now();
|
||||
let entry = AuditLogEntry {
|
||||
|
|
@ -904,7 +992,7 @@ async fn handle_create_audit(State(state): State<ApiState>, Json(req): Json<Crea
|
|||
logs.truncate(100);
|
||||
}
|
||||
|
||||
ApiResponse::success(true)
|
||||
(StatusCode::OK, ApiResponse::success(true))
|
||||
}
|
||||
|
||||
// ── Bulk keys & Router Rules ─────────────────────────────────────────────────
|
||||
|
|
@ -1006,10 +1094,16 @@ async fn handle_put_rules(
|
|||
(StatusCode::OK, ApiResponse::success(true))
|
||||
}
|
||||
|
||||
async fn handle_clear_audit(State(state): State<ApiState>) -> impl IntoResponse {
|
||||
let mut logs = state.audit_logs.write().unwrap();
|
||||
async fn handle_clear_audit(
|
||||
State(state): State<ApiState>,
|
||||
headers: axum::http::HeaderMap,
|
||||
) -> impl IntoResponse {
|
||||
if !check_token(&state, &headers) {
|
||||
return api_unauthorized::<()>();
|
||||
}
|
||||
let mut logs = state.audit_logs.write().unwrap_or_else(|e| e.into_inner());
|
||||
logs.clear();
|
||||
ApiResponse::success(())
|
||||
(StatusCode::OK, ApiResponse::success(()))
|
||||
}
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -276,6 +276,18 @@ impl DnsServer {
|
|||
///
|
||||
/// Клиент может явно указать `<server_ip>:<local_port>` как DNS-сервер
|
||||
/// в настройках — тогда все DNS-запросы туннелируются и резолвятся здесь.
|
||||
///
|
||||
/// SECURITY: this socket is bound on 0.0.0.0, reachable directly from the
|
||||
/// public internet with no authentication (unlike the main OSTP port,
|
||||
/// there is no Noise handshake gating it). Answering every UDP datagram
|
||||
/// by resolving and replying to its (unverified, spoofable) source
|
||||
/// address is a textbook DNS reflection/amplification primitive: an
|
||||
/// attacker spoofing a victim's IP as the query source turns this server
|
||||
/// into a free amplifier against that victim. There is currently no
|
||||
/// caller for this function anywhere in the codebase, but the rate
|
||||
/// limiter below exists so that connecting it later doesn't silently
|
||||
/// reintroduce that risk - it bounds how much amplification bandwidth
|
||||
/// this listener can ever contribute, regardless of query volume.
|
||||
pub async fn run_local_udp_listener(self: Arc<Self>) {
|
||||
let port = self.config.read().await.local_port;
|
||||
let bind_addr = format!("0.0.0.0:{port}");
|
||||
|
|
@ -289,10 +301,30 @@ impl DnsServer {
|
|||
};
|
||||
tracing::info!("Built-in DNS server listening on UDP {bind_addr}");
|
||||
|
||||
// Global token bucket capping total replies/sec this listener will
|
||||
// ever send. Deliberately global (not per-source-IP): per-IP limiting
|
||||
// does nothing against a reflection attack, since the attacker never
|
||||
// sees the responses and can spread queries across arbitrarily many
|
||||
// spoofed sources anyway. A global cap bounds this server's total
|
||||
// contribution to any attack regardless of how the queries are
|
||||
// distributed.
|
||||
const MAX_REPLIES_PER_SEC: f64 = 100.0;
|
||||
let mut tokens: f64 = MAX_REPLIES_PER_SEC;
|
||||
let mut last_refill = tokio::time::Instant::now();
|
||||
|
||||
let mut buf = vec![0u8; 4096];
|
||||
loop {
|
||||
match socket.recv_from(&mut buf).await {
|
||||
Ok((n, peer)) => {
|
||||
let now = tokio::time::Instant::now();
|
||||
tokens = (tokens + now.duration_since(last_refill).as_secs_f64() * MAX_REPLIES_PER_SEC)
|
||||
.min(MAX_REPLIES_PER_SEC);
|
||||
last_refill = now;
|
||||
if tokens < 1.0 {
|
||||
continue; // over budget: drop silently, no reply sent
|
||||
}
|
||||
tokens -= 1.0;
|
||||
|
||||
let query = buf[..n].to_vec();
|
||||
let srv = self.clone();
|
||||
let sock = socket.clone();
|
||||
|
|
|
|||
|
|
@ -21,3 +21,4 @@ tracing-subscriber = { version = "0.3", features = ["env-filter"] }
|
|||
ostp-core = { path = "../ostp-core" }
|
||||
colored = "2.1"
|
||||
rlimit = "0.11.0"
|
||||
sha2.workspace = true
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@ use clap::Parser;
|
|||
use std::fs;
|
||||
use std::path::PathBuf;
|
||||
use colored::Colorize;
|
||||
use sha2::Digest;
|
||||
|
||||
#[derive(Parser, Debug)]
|
||||
#[command(author, version, about = "OSTP Core - Ospab Stealth Transport Protocol", long_about = None)]
|
||||
|
|
@ -720,24 +721,11 @@ fn run_setup_wizard(config_path: &std::path::Path) -> Result<()> {
|
|||
}) as char
|
||||
}).collect();
|
||||
let password = wizard_prompt("Admin password (blank for random)", &rand_pass);
|
||||
let pass_hash = {
|
||||
use std::fmt::Write as _;
|
||||
let mut hash = String::new();
|
||||
let digest: [u8; 32] = {
|
||||
use std::collections::hash_map::DefaultHasher;
|
||||
use std::hash::{Hash, Hasher};
|
||||
// Panel password hashing. sha2 is not a direct dep of ostp/Cargo.toml,
|
||||
// so we use std's hasher as a placeholder digest here.
|
||||
let mut h = DefaultHasher::new();
|
||||
password.hash(&mut h);
|
||||
let v = h.finish();
|
||||
let mut out = [0u8; 32];
|
||||
out[..8].copy_from_slice(&v.to_be_bytes());
|
||||
out
|
||||
};
|
||||
for b in digest { let _ = write!(hash, "{:02x}", b); }
|
||||
hash
|
||||
};
|
||||
// Must match api.rs's handle_login exactly (format!("{:x}", Sha256::digest(..))) -
|
||||
// this used to be a DefaultHasher (SipHash) placeholder that produced a
|
||||
// differently-shaped digest, so a password set up through this wizard could
|
||||
// never actually log into the panel it just configured.
|
||||
let pass_hash = format!("{:x}", sha2::Sha256::digest(password.as_bytes()));
|
||||
|
||||
wizard_step(4, TOTAL, "Saving configuration");
|
||||
let panel_bind = format!("0.0.0.0:{}", panel_port);
|
||||
|
|
|
|||
Loading…
Reference in New Issue