Carrier negotiation + WebSocket lifecycle

Co-Authored-By: brekotis <93345790+brekotis@users.noreply.github.com>
This commit is contained in:
Alexey
2026-08-27 07:38:42 +03:00
parent 80a2737eed
commit c75cf5cc9d
28 changed files with 1411 additions and 1320 deletions
+12 -3
View File
@@ -110,19 +110,28 @@ impl WebProcessRuntime {
})
}
/// Resolves non-secret bootstrap trace identity without exposing its credential.
/// Resolves bootstrap trace identity and the frozen live-session body timeout.
pub(crate) fn bootstrap_trace_identity(
&self,
hash: TokenHash,
host: &str,
) -> Option<(u64, Arc<WebRuntimeProfile>)> {
) -> Option<(u64, Arc<WebRuntimeProfile>, Option<Duration>)> {
let now = Instant::now();
self.state
.lock()
.bootstraps
.get(&hash)
.filter(|entry| entry.profile.host == host && now <= entry.expires_at)
.map(|entry| (entry.trace_session_id, Arc::clone(&entry.profile)))
.map(|entry| {
(
entry.trace_session_id,
Arc::clone(&entry.profile),
entry
.session
.as_ref()
.map(|session| Duration::from_secs(session.timeouts().body_secs)),
)
})
}
/// Resolves an authenticated session token.
-346
View File
@@ -1,346 +0,0 @@
/// Telemt Carrier Selection and Failure Dampening - Copyright 2077
/// anhand des Kundenverhaltens Rückschlüsse gegen DSGVO ziehen...?!
use std::collections::HashMap;
use std::net::IpAddr;
use std::time::{Duration, Instant};
use sha2::{Digest, Sha256};
use super::negotiation::{CarrierClientClass, CarrierLearningContext};
use super::ProfileKey;
use crate::config::WebCarrier;
const PROFILE_WEIGHT: i16 = 4;
const USER_AGENT_WEIGHT: i16 = 4;
const IP_WEIGHT: i16 = 1;
const SCORE_MIN: i8 = -8;
const SCORE_MAX: i8 = 8;
const PROFILE_MIN_OUTCOMES: u8 = 8;
const PROFILE_MIN_COHORTS: usize = 4;
const COHORT_CONTEXT: &[u8] = b"telemt-web-carrier-cohort-v1\0";
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
enum EvidenceKey {
Profile(ProfileKey),
UserAgent(ProfileKey, CarrierClientClass, [u8; 32]),
Ip(ProfileKey, IpAddr),
}
struct Evidence {
created_at: Instant,
lifetime: Duration,
scores: [i8; 4],
outcomes: u8,
cohorts: [Option<[u8; 32]>; PROFILE_MIN_COHORTS],
}
impl Evidence {
fn new(created_at: Instant, lifetime: Duration) -> Self {
Self {
created_at,
lifetime,
scores: [0; 4],
outcomes: 0,
cohorts: [None; PROFILE_MIN_COHORTS],
}
}
fn update(&mut self, carrier: WebCarrier, delta: i8, cohort: Option<[u8; 32]>) {
let score = &mut self.scores[carrier.index()];
*score = score.saturating_add(delta).clamp(SCORE_MIN, SCORE_MAX);
self.outcomes = self.outcomes.saturating_add(1).min(PROFILE_MIN_OUTCOMES);
if let Some(cohort) = cohort
&& !self.cohorts.contains(&Some(cohort))
&& let Some(slot) = self.cohorts.iter_mut().find(|slot| slot.is_none())
{
*slot = Some(cohort);
}
}
}
/// Process-local bounded fixed-window carrier evidence store.
pub(super) struct CarrierLearning {
entries: HashMap<EvidenceKey, Evidence>,
capacity: usize,
}
impl CarrierLearning {
/// Creates an empty store under the restart-owned capacity ceiling.
pub(super) fn new(capacity: usize) -> Self {
Self {
entries: HashMap::with_capacity(capacity),
capacity,
}
}
/// Ranks supported configured candidates using only unexpired evidence.
pub(super) fn rank(
&mut self,
now: Instant,
configured: &[WebCarrier],
request: super::CarrierRequest,
profile_key: ProfileKey,
client_ip: IpAddr,
) -> (Vec<WebCarrier>, [i16; 4]) {
self.prune(now);
let mut scores = [0i16; 4];
let profile = self.entries.get(&EvidenceKey::Profile(profile_key));
let profile_ready = profile.is_some_and(|entry| {
entry.outcomes >= PROFILE_MIN_OUTCOMES
&& entry.cohorts.iter().flatten().count() >= PROFILE_MIN_COHORTS
});
let user_agent = self.entries.get(&EvidenceKey::UserAgent(
profile_key,
request.class(),
request.user_agent_hash(),
));
let ip = self.entries.get(&EvidenceKey::Ip(profile_key, client_ip));
for carrier in WebCarrier::ALL {
let index = carrier.index();
if profile_ready {
scores[index] += i16::from(profile.map_or(0, |entry| entry.scores[index]))
* PROFILE_WEIGHT;
}
scores[index] += i16::from(user_agent.map_or(0, |entry| entry.scores[index]))
* USER_AGENT_WEIGHT;
scores[index] +=
i16::from(ip.map_or(0, |entry| entry.scores[index])) * IP_WEIGHT;
}
let mut ranked = configured
.iter()
.copied()
.filter(|carrier| request.supports(*carrier))
.collect::<Vec<_>>();
ranked.sort_by_key(|carrier| std::cmp::Reverse(scores[carrier.index()]));
(ranked, scores)
}
/// Records one committed success or one server-accepted supersession failure.
pub(super) fn record(
&mut self,
now: Instant,
lifetime: Duration,
context: CarrierLearningContext,
carrier: WebCarrier,
success: bool,
) {
self.prune(now);
let delta = if success { 1 } else { -1 };
let cohort = cohort_hash(context);
self.update(
EvidenceKey::Profile(context.profile_key),
now,
lifetime,
carrier,
delta,
Some(cohort),
);
self.update(
EvidenceKey::UserAgent(
context.profile_key,
context.class,
context.user_agent_hash,
),
now,
lifetime,
carrier,
delta,
None,
);
self.update(
EvidenceKey::Ip(context.profile_key, context.client_ip),
now,
lifetime,
carrier,
delta,
None,
);
}
/// Removes fixed-window entries after their creation-time expiry.
pub(super) fn prune(&mut self, now: Instant) {
self.entries.retain(|_, entry| {
now.saturating_duration_since(entry.created_at) <= entry.lifetime
});
}
fn update(
&mut self,
key: EvidenceKey,
now: Instant,
lifetime: Duration,
carrier: WebCarrier,
delta: i8,
cohort: Option<[u8; 32]>,
) {
if !self.entries.contains_key(&key) && self.entries.len() >= self.capacity {
let oldest = self
.entries
.iter()
.min_by_key(|(_, entry)| entry.created_at)
.map(|(key, _)| *key);
if let Some(oldest) = oldest {
self.entries.remove(&oldest);
}
}
self.entries
.entry(key)
.or_insert_with(|| Evidence::new(now, lifetime))
.update(carrier, delta, cohort);
}
}
fn cohort_hash(context: CarrierLearningContext) -> [u8; 32] {
let mut digest = Sha256::new();
digest.update(COHORT_CONTEXT);
digest.update(context.profile_key);
digest.update([match context.class {
CarrierClientClass::Legacy => 0,
CarrierClientClass::Bridge => 1,
CarrierClientClass::BrowserHint => 2,
}]);
digest.update(context.user_agent_hash);
match context.client_ip {
IpAddr::V4(address) => {
digest.update([4]);
digest.update(address.octets());
}
IpAddr::V6(address) => {
digest.update([6]);
digest.update(address.octets());
}
}
digest.finalize().into()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::web::manager::{CarrierCapabilities, CarrierRequest};
fn request(hash: u8) -> CarrierRequest {
CarrierRequest::automatic(
CarrierClientClass::Bridge,
CarrierCapabilities::all(),
1,
None,
[hash; 32],
)
}
#[test]
fn evidence_is_bounded_and_expires_without_sliding() {
let start = Instant::now();
let mut learning = CarrierLearning::new(3);
let context = CarrierLearningContext {
profile_key: [1; 32],
client_ip: "192.0.2.1".parse().unwrap(),
class: CarrierClientClass::Bridge,
user_agent_hash: [2; 32],
};
learning.record(
start,
Duration::from_secs(10),
context,
WebCarrier::Websocket,
true,
);
assert_eq!(learning.entries.len(), 3);
learning.record(
start + Duration::from_secs(5),
Duration::from_secs(10),
context,
WebCarrier::Websocket,
true,
);
learning.prune(start + Duration::from_secs(11));
assert!(learning.entries.is_empty());
}
#[test]
fn user_agent_and_ip_evidence_rank_stably() {
let now = Instant::now();
let mut learning = CarrierLearning::new(16);
let context = CarrierLearningContext {
profile_key: [1; 32],
client_ip: "192.0.2.1".parse().unwrap(),
class: CarrierClientClass::Bridge,
user_agent_hash: [2; 32],
};
learning.record(
now,
Duration::from_secs(10),
context,
WebCarrier::Websocket,
true,
);
let (ranked, scores) = learning.rank(
now,
&[WebCarrier::Https, WebCarrier::Websocket],
request(2),
context.profile_key,
context.client_ip,
);
assert_eq!(ranked, [WebCarrier::Websocket, WebCarrier::Https]);
assert_eq!(scores[WebCarrier::Websocket.index()], 5);
}
#[test]
fn profile_evidence_requires_outcome_and_cohort_thresholds() {
let now = Instant::now();
let profile_key = [1; 32];
let configured = [WebCarrier::Https, WebCarrier::Websocket];
let unrelated_ip = "198.51.100.10".parse().unwrap();
let mut learning = CarrierLearning::new(64);
for cohort in 1..=3u8 {
let context = CarrierLearningContext {
profile_key,
client_ip: IpAddr::V4(std::net::Ipv4Addr::new(192, 0, 2, cohort)),
class: CarrierClientClass::Bridge,
user_agent_hash: [cohort; 32],
};
for _ in 0..2 {
learning.record(
now,
Duration::from_secs(10),
context,
WebCarrier::Websocket,
true,
);
}
}
let (ranked, _) = learning.rank(
now,
&configured,
request(99),
profile_key,
unrelated_ip,
);
assert_eq!(ranked, configured);
let fourth = CarrierLearningContext {
profile_key,
client_ip: "192.0.2.4".parse().unwrap(),
class: CarrierClientClass::Bridge,
user_agent_hash: [4; 32],
};
for _ in 0..2 {
learning.record(
now,
Duration::from_secs(10),
fourth,
WebCarrier::Websocket,
true,
);
}
let (ranked, scores) = learning.rank(
now,
&configured,
request(99),
profile_key,
unrelated_ip,
);
assert_eq!(ranked, [WebCarrier::Websocket, WebCarrier::Https]);
assert_eq!(scores[WebCarrier::Websocket.index()], 32);
}
}
+2 -8
View File
@@ -17,6 +17,7 @@ impl WebProcessRuntime {
client_ip: IpAddr,
profile_key: ProfileKey,
profile_host: &str,
closed_token_lifetime: Duration,
) {
let mut state = self.state.lock();
if state.sessions.remove(&hash).is_none() {
@@ -28,14 +29,7 @@ impl WebProcessRuntime {
&mut state,
hash,
profile_host,
Duration::from_secs(
self.active_runtime
.load()
.config()
.web
.timeouts
.bootstrap_lifetime_secs,
),
closed_token_lifetime,
self.limits.max_sessions_global.saturating_mul(16),
);
let bootstrap_hashes = state
+10
View File
@@ -77,6 +77,11 @@ impl CarrierCapabilities {
Self(0b1111)
}
/// Returns the current server-authoritative native iOS capability ceiling.
pub(crate) const fn ios() -> Self {
Self(1 << WebCarrier::Https.index())
}
/// Builds a set from a validated bit representation.
pub(crate) const fn from_bits(bits: u8) -> Option<Self> {
if bits != 0 && bits & !0b1111 == 0 {
@@ -90,6 +95,11 @@ impl CarrierCapabilities {
pub(crate) const fn contains(self, carrier: WebCarrier) -> bool {
self.0 & (1 << carrier.index()) != 0
}
/// Intersects declared capabilities with an authoritative server ceiling.
pub(crate) const fn intersection(self, ceiling: Self) -> Option<Self> {
Self::from_bits(self.0 & ceiling.0)
}
}
/// Immutable metadata attached to one session-creation attempt.
+3 -174
View File
@@ -376,178 +376,7 @@ impl WebProcessRuntime {
);
Ok(result)
}
fn replace_session(
self: &Arc<Self>,
bootstrap_hash: TokenHash,
client_ip: IpAddr,
replacement: Replacement,
) -> std::result::Result<CreateResult, ManagerError> {
if !replacement.old_session.begin_carrier_supersede() {
let committed = replacement.old_session.is_carrier_committed();
self.cancel_replacement(bootstrap_hash, &replacement.old_session);
return Err(if committed {
ManagerError::Committed
} else {
ManagerError::Closed
});
}
let generation = self.active_generation();
let config = generation.config();
let now = Instant::now();
let mut state = self.state.lock();
remove_expired_locked(&mut state, now);
let valid = state.bootstraps.get(&bootstrap_hash).is_some_and(|entry| {
entry.carrier_transitioning
&& entry.carrier_phase == CarrierChainPhase::Provisional
&& !entry.close_requested
&& entry.carrier_attempt.saturating_add(1) == replacement.attempt
&& now < replacement.carrier_deadline_at
&& entry
.session
.as_ref()
.is_some_and(|session| Arc::ptr_eq(session, &replacement.old_session))
}) && state
.sessions
.get(&replacement.old_session.token_hash())
.is_some_and(|session| Arc::ptr_eq(session, &replacement.old_session));
if !valid
|| state.closed
|| !config.web.enabled
|| !generation
.proxy_shared
.is_user_enabled(&replacement.profile.user)
{
drop(state);
self.cancel_replacement(bootstrap_hash, &replacement.old_session);
return Err(ManagerError::Closed);
}
let Some((session_token, session_hash)) = new_unique_token(&generation, &state) else {
self.limit_hits.fetch_add(1, Ordering::Relaxed);
drop(state);
self.cancel_replacement(bootstrap_hash, &replacement.old_session);
return Err(ManagerError::Limit);
};
let learning_context = (replacement.profile.carrier_learning
&& replacement.learning_epoch != 0)
.then_some(CarrierLearningContext {
profile_key: replacement.profile_key,
client_ip,
class: replacement.request.class(),
user_agent_hash: replacement.request.user_agent_hash(),
epoch: replacement.learning_epoch,
ip_learning_eligible: replacement.ip_learning_eligible,
});
let session = WebSession::new(
Arc::downgrade(self),
session_hash,
client_ip,
replacement.trace_session_id,
Arc::clone(&replacement.profile),
replacement.profile_key,
replacement.carrier,
replacement.attempt,
bootstrap_hash,
Some(replacement.carrier_deadline_at),
replacement.request.class(),
learning_context,
true,
self.limits.clone(),
replacement.old_session.timeouts().clone(),
);
let Some(supersede) = replacement.old_session.prepare_carrier_supersede() else {
drop(state);
self.cancel_replacement(bootstrap_hash, &replacement.old_session);
session.close();
return Err(ManagerError::Closed);
};
let old_hash = replacement.old_session.token_hash();
state.sessions.remove(&old_hash);
remember_closed_token_locked(
&mut state,
old_hash,
&replacement.profile.host,
Duration::from_secs(config.web.timeouts.bootstrap_lifetime_secs),
self.limits.max_sessions_global.saturating_mul(16),
);
state.sessions.insert(session_hash, Arc::clone(&session));
let entry = state
.bootstraps
.get_mut(&bootstrap_hash)
.ok_or(ManagerError::Authentication)?;
entry.session_token = Zeroizing::new(session_token.clone());
entry.session = Some(Arc::clone(&session));
entry.carrier_request = Some(replacement.request);
entry.carrier_attempt = replacement.attempt;
entry.carrier_transitioning = false;
entry.carrier_phase = CarrierChainPhase::Provisional;
if let Some(slot) = entry
.carrier_failures
.get_mut(usize::from(replacement.attempt.saturating_sub(2)))
{
*slot = Some(replacement.old_session.carrier());
}
self.sessions_created.fetch_add(1, Ordering::Relaxed);
self.sessions_closed.fetch_add(1, Ordering::Relaxed);
let result = CreateResult {
token: session_token,
carrier: replacement.carrier,
attempt: Some(replacement.attempt),
candidate_count: Some(u8::try_from(entry.carrier_candidates.len()).unwrap_or(4)),
deadline_secs: Some(entry.profile.carrier_negotiation_deadlines_secs[3]),
carrier_state: Some(CarrierChainPhase::Provisional.as_str()),
};
let identity = session.trace_identity();
let old_identity = replacement.old_session.trace_identity();
drop(state);
supersede.finish();
self.trace.record_carrier_lifecycle(
client_ip,
old_identity.clone(),
TraceLifecycleEvent::CarrierFailed,
replacement.request.class().as_str(),
replacement.old_session.carrier(),
replacement.attempt - 1,
replacement.scores,
replacement
.request
.failure()
.map(|failure| failure.as_str()),
);
self.trace.record_carrier_lifecycle(
client_ip,
old_identity,
TraceLifecycleEvent::CarrierSuperseded,
replacement.request.class().as_str(),
replacement.old_session.carrier(),
replacement.attempt - 1,
replacement.scores,
replacement
.request
.failure()
.map(|failure| failure.as_str()),
);
self.trace.record_carrier_lifecycle(
client_ip,
identity.clone(),
TraceLifecycleEvent::CarrierSelected,
replacement.request.class().as_str(),
replacement.carrier,
replacement.attempt,
replacement.scores,
None,
);
self.trace.record_lifecycle(
None,
Some(client_ip),
identity,
TraceLifecycleEvent::SessionCreated,
None,
replacement
.request
.failure()
.map(|failure| failure.as_str()),
);
Ok(result)
}
}
// Atomic pre-commit carrier replacement and frozen-policy transfer.
mod replacement;
@@ -0,0 +1,177 @@
use super::*;
impl WebProcessRuntime {
pub(super) fn replace_session(
self: &Arc<Self>,
bootstrap_hash: TokenHash,
client_ip: IpAddr,
replacement: Replacement,
) -> std::result::Result<CreateResult, ManagerError> {
if !replacement.old_session.begin_carrier_supersede() {
let committed = replacement.old_session.is_carrier_committed();
self.cancel_replacement(bootstrap_hash, &replacement.old_session);
return Err(if committed {
ManagerError::Committed
} else {
ManagerError::Closed
});
}
let generation = self.active_generation();
let config = generation.config();
let now = Instant::now();
let mut state = self.state.lock();
remove_expired_locked(&mut state, now);
let valid = state.bootstraps.get(&bootstrap_hash).is_some_and(|entry| {
entry.carrier_transitioning
&& entry.carrier_phase == CarrierChainPhase::Provisional
&& !entry.close_requested
&& entry.carrier_attempt.saturating_add(1) == replacement.attempt
&& now < replacement.carrier_deadline_at
&& entry
.session
.as_ref()
.is_some_and(|session| Arc::ptr_eq(session, &replacement.old_session))
}) && state
.sessions
.get(&replacement.old_session.token_hash())
.is_some_and(|session| Arc::ptr_eq(session, &replacement.old_session));
if !valid
|| state.closed
|| !config.web.enabled
|| !generation
.proxy_shared
.is_user_enabled(&replacement.profile.user)
{
drop(state);
self.cancel_replacement(bootstrap_hash, &replacement.old_session);
return Err(ManagerError::Closed);
}
let Some((session_token, session_hash)) = new_unique_token(&generation, &state) else {
self.limit_hits.fetch_add(1, Ordering::Relaxed);
drop(state);
self.cancel_replacement(bootstrap_hash, &replacement.old_session);
return Err(ManagerError::Limit);
};
let learning_context = (replacement.profile.carrier_learning
&& replacement.learning_epoch != 0)
.then_some(CarrierLearningContext {
profile_key: replacement.profile_key,
client_ip,
class: replacement.request.class(),
user_agent_hash: replacement.request.user_agent_hash(),
epoch: replacement.learning_epoch,
ip_learning_eligible: replacement.ip_learning_eligible,
});
let session = WebSession::new(
Arc::downgrade(self),
session_hash,
client_ip,
replacement.trace_session_id,
Arc::clone(&replacement.profile),
replacement.profile_key,
replacement.carrier,
replacement.attempt,
bootstrap_hash,
Some(replacement.carrier_deadline_at),
replacement.request.class(),
learning_context,
true,
self.limits.clone(),
replacement.old_session.timeouts().clone(),
);
let Some(supersede) = replacement.old_session.prepare_carrier_supersede() else {
drop(state);
self.cancel_replacement(bootstrap_hash, &replacement.old_session);
session.close();
return Err(ManagerError::Closed);
};
let old_hash = replacement.old_session.token_hash();
state.sessions.remove(&old_hash);
remember_closed_token_locked(
&mut state,
old_hash,
&replacement.profile.host,
Duration::from_secs(replacement.old_session.timeouts().bootstrap_lifetime_secs),
self.limits.max_sessions_global.saturating_mul(16),
);
state.sessions.insert(session_hash, Arc::clone(&session));
let entry = state
.bootstraps
.get_mut(&bootstrap_hash)
.ok_or(ManagerError::Authentication)?;
entry.session_token = Zeroizing::new(session_token.clone());
entry.session = Some(Arc::clone(&session));
entry.carrier_request = Some(replacement.request);
entry.carrier_attempt = replacement.attempt;
entry.carrier_transitioning = false;
entry.carrier_phase = CarrierChainPhase::Provisional;
if let Some(slot) = entry
.carrier_failures
.get_mut(usize::from(replacement.attempt.saturating_sub(2)))
{
*slot = Some(replacement.old_session.carrier());
}
self.sessions_created.fetch_add(1, Ordering::Relaxed);
self.sessions_closed.fetch_add(1, Ordering::Relaxed);
let result = CreateResult {
token: session_token,
carrier: replacement.carrier,
attempt: Some(replacement.attempt),
candidate_count: Some(u8::try_from(entry.carrier_candidates.len()).unwrap_or(4)),
deadline_secs: Some(entry.profile.carrier_negotiation_deadlines_secs[3]),
carrier_state: Some(CarrierChainPhase::Provisional.as_str()),
};
let identity = session.trace_identity();
let old_identity = replacement.old_session.trace_identity();
drop(state);
supersede.finish();
self.trace.record_carrier_lifecycle(
client_ip,
old_identity.clone(),
TraceLifecycleEvent::CarrierFailed,
replacement.request.class().as_str(),
replacement.old_session.carrier(),
replacement.attempt - 1,
replacement.scores,
replacement
.request
.failure()
.map(|failure| failure.as_str()),
);
self.trace.record_carrier_lifecycle(
client_ip,
old_identity,
TraceLifecycleEvent::CarrierSuperseded,
replacement.request.class().as_str(),
replacement.old_session.carrier(),
replacement.attempt - 1,
replacement.scores,
replacement
.request
.failure()
.map(|failure| failure.as_str()),
);
self.trace.record_carrier_lifecycle(
client_ip,
identity.clone(),
TraceLifecycleEvent::CarrierSelected,
replacement.request.class().as_str(),
replacement.carrier,
replacement.attempt,
replacement.scores,
None,
);
self.trace.record_lifecycle(
None,
Some(client_ip),
identity,
TraceLifecycleEvent::SessionCreated,
None,
replacement
.request
.failure()
.map(|failure| failure.as_str()),
);
Ok(result)
}
}