Skip to content

Commit a8559be

Browse files
author
pseusys
committed
fixes - fixed again
1 parent b8a4107 commit a8559be

6 files changed

Lines changed: 43 additions & 68 deletions

File tree

typhoon/src/flow/config.rs

Lines changed: 18 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -187,16 +187,30 @@ impl FakeHeaderConfig {
187187
/// Create a random header configuration drawn from the default probability distributions.
188188
///
189189
/// Includes a header with probability `FAKE_HEADER_PROBABILITY`; if included, a random number
190-
/// of U8-random fields are packed to fill a length sampled from
191-
/// `[FAKE_HEADER_LENGTH_MIN, FAKE_HEADER_LENGTH_MAX]`; otherwise returns an empty config.
190+
/// of fields are packed to fill a length sampled from `[FAKE_HEADER_LENGTH_MIN, FAKE_HEADER_LENGTH_MAX]`.
191+
/// Each field is independently assigned one of the five `FieldType` variants with equal probability.
192192
pub fn random<AE: AsyncExecutor>(settings: &Settings<AE>) -> Self {
193193
let mut rng = get_rng();
194194
let header_prob = settings.get(&keys::FAKE_HEADER_PROBABILITY);
195195
if rng.r#gen::<f64>() < header_prob {
196196
let min_len = settings.get(&keys::FAKE_HEADER_LENGTH_MIN) as usize;
197197
let max_len = settings.get(&keys::FAKE_HEADER_LENGTH_MAX) as usize;
198198
let len = if min_len >= max_len { max_len } else { rng.gen_range(min_len..=max_len) };
199-
let fields = (0..len).map(|_| FieldTypeHolder::U8(FieldType::Random)).collect();
199+
let fields = (0..len)
200+
.map(|_| {
201+
FieldTypeHolder::U8(match rng.gen_range(0u8..5) {
202+
0 => FieldType::Random,
203+
1 => FieldType::Constant { value: rng.r#gen::<u8>() },
204+
2 => FieldType::Volatile { value: rng.r#gen::<u8>(), change_probability: rng.gen_range(0.01..=0.20) },
205+
3 => {
206+
let switch_timeout = rng.gen_range(1_000u64..=30_000);
207+
FieldType::Switching { value: rng.r#gen::<u8>(), next_switch: unix_timestamp_ms() + switch_timeout as u128, switch_timeout }
208+
}
209+
4 => FieldType::Incremental { value: rng.r#gen::<u8>() },
210+
_ => unreachable!(),
211+
})
212+
})
213+
.collect();
200214
Self::new(fields)
201215
} else {
202216
Self::new(vec![])
@@ -276,46 +290,7 @@ impl FlowConfig {
276290
/// `[FAKE_BODY_LENGTH_MIN, mtu]`.
277291
pub fn random<AE: AsyncExecutor>(settings: &Settings<AE>) -> Self {
278292
let mut rng = get_rng();
279-
280-
let header_prob = settings.get(&keys::FAKE_HEADER_PROBABILITY);
281-
let fake_header_mode = if rng.r#gen::<f64>() < header_prob {
282-
let min_len = settings.get(&keys::FAKE_HEADER_LENGTH_MIN) as usize;
283-
let max_len = settings.get(&keys::FAKE_HEADER_LENGTH_MAX) as usize;
284-
let len = if min_len >= max_len {
285-
max_len
286-
} else {
287-
rng.gen_range(min_len..=max_len)
288-
};
289-
let fields = (0..len)
290-
.map(|_| {
291-
FieldTypeHolder::U8(match rng.gen_range(0u8..5) {
292-
0 => FieldType::Random,
293-
1 => FieldType::Constant {
294-
value: rng.r#gen::<u8>(),
295-
},
296-
2 => FieldType::Volatile {
297-
value: rng.r#gen::<u8>(),
298-
change_probability: rng.gen_range(0.01..=0.20),
299-
},
300-
3 => {
301-
let switch_timeout = rng.gen_range(1_000u64..=30_000);
302-
FieldType::Switching {
303-
value: rng.r#gen::<u8>(),
304-
next_switch: unix_timestamp_ms() + switch_timeout as u128,
305-
switch_timeout,
306-
}
307-
}
308-
4 => FieldType::Incremental {
309-
value: rng.r#gen::<u8>(),
310-
},
311-
_ => unreachable!(),
312-
})
313-
})
314-
.collect();
315-
FakeHeaderConfig::new(fields)
316-
} else {
317-
FakeHeaderConfig::new(vec![])
318-
};
293+
let fake_header_mode = FakeHeaderConfig::random(settings);
319294

320295
let min_len = settings.get(&keys::FAKE_BODY_LENGTH_MIN) as usize;
321296
let max_len = settings.get(&keys::FAKE_BODY_LENGTH_MAX) as usize;

typhoon/src/flow/decoy/common.rs

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -516,7 +516,7 @@ where
516516
break;
517517
};
518518

519-
let (packet, body_length, should_rep) = {
519+
let (packet, body_length, should_rep, settings) = {
520520
let mut guard = state.write().await;
521521
let length = guard.pending_maintenance_length;
522522

@@ -527,18 +527,18 @@ where
527527

528528
let packet = guard.create_decoy_packet(length, true);
529529
let should_rep = guard.should_replicate(true);
530-
(packet, length, should_rep)
530+
let settings = Arc::clone(&guard.settings);
531+
(packet, length, should_rep, settings)
531532
};
532533

533-
// Clone before consuming: Arc refcount bump (zero-copy) for optional replication.
534-
let body_bytes = should_rep.then(|| packet.slice_end(body_length).to_vec());
534+
let body_buf = should_rep.then(|| settings.pool().allocate_precise_from_slice_with_capacity(packet.slice_end(body_length), 0, 0));
535535

536536
debug!("Maintenance: generated packet (len={body_length})");
537537

538538
if let Err(err) = manager_arc.send_decoy_packet(packet).await {
539539
warn!("Maintenance: failed to send: {err:?}");
540-
} else if let Some(bytes) = body_bytes {
541-
try_replicate(&state, &manager, true, bytes).await;
540+
} else if let Some(body) = body_buf {
541+
try_replicate(&state, &manager, true, body).await;
542542
}
543543

544544
{
@@ -550,7 +550,7 @@ where
550550

551551
/// Attempt replication of a decoy packet. If replication mode applies, spawns a cascading
552552
/// task that re-sends the packet body with diminishing probability.
553-
pub(super) async fn try_replicate<T, AE>(state: &Arc<RwLock<DecoyState<T, AE>>>, manager: &Weak<dyn DecoyFlowSender>, is_maintenance: bool, body_bytes: Vec<u8>)
553+
pub(super) async fn try_replicate<T, AE>(state: &Arc<RwLock<DecoyState<T, AE>>>, manager: &Weak<dyn DecoyFlowSender>, is_maintenance: bool, body: DynamicByteBuffer)
554554
where
555555
T: IdentityType + Clone + 'static,
556556
AE: AsyncExecutor + 'static,
@@ -582,10 +582,10 @@ where
582582

583583
let packet = {
584584
let mut guard = state_clone.write().await;
585-
if !guard.try_spend_budget(body_bytes.len()) {
585+
if !guard.try_spend_budget(body.slice().len()) {
586586
break;
587587
}
588-
guard.create_replica_packet(&body_bytes, is_maintenance)
588+
guard.create_replica_packet(body.slice(), is_maintenance)
589589
};
590590

591591
if manager_arc.send_decoy_packet(packet).await.is_err() {

typhoon/src/flow/decoy/heavy.rs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -75,14 +75,14 @@ impl<T: IdentityType + Clone, AE: AsyncExecutor> HeavyDecoyProvider<T, AE> {
7575
if state_guard.try_spend_budget(decoy_length) {
7676
let decoy_packet = state_guard.create_decoy_packet(decoy_length, false);
7777
let should_rep = state_guard.should_replicate(false);
78+
let settings = Arc::clone(&state_guard.settings);
7879
drop(state_guard);
7980

80-
// Clone before consuming: Arc refcount bump (zero-copy) for optional replication.
81-
let body_bytes = should_rep.then(|| decoy_packet.slice_end(decoy_length).to_vec());
81+
let body_buf = should_rep.then(|| settings.pool().allocate_precise_from_slice_with_capacity(decoy_packet.slice_end(decoy_length), 0, 0));
8282
if let Err(err) = manager_arc.send_decoy_packet(decoy_packet).await {
8383
warn!("HeavyDecoyProvider: failed to send decoy packet: {err:?}");
84-
} else if let Some(bytes) = body_bytes {
85-
try_replicate(&state, &manager, false, bytes).await;
84+
} else if let Some(body) = body_buf {
85+
try_replicate(&state, &manager, false, body).await;
8686
}
8787
}
8888
}

typhoon/src/flow/decoy/noisy.rs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -73,14 +73,14 @@ impl<T: IdentityType + Clone, AE: AsyncExecutor> NoisyDecoyProvider<T, AE> {
7373
if state_guard.try_spend_budget(decoy_length) {
7474
let decoy_packet = state_guard.create_decoy_packet(decoy_length, false);
7575
let should_rep = state_guard.should_replicate(false);
76+
let settings = Arc::clone(&state_guard.settings);
7677
drop(state_guard);
7778

78-
// Clone before consuming: Arc refcount bump (zero-copy) for optional replication.
79-
let body_bytes = should_rep.then(|| decoy_packet.slice_end(decoy_length).to_vec());
79+
let body_buf = should_rep.then(|| settings.pool().allocate_precise_from_slice_with_capacity(decoy_packet.slice_end(decoy_length), 0, 0));
8080
if let Err(err) = manager_arc.send_decoy_packet(decoy_packet).await {
8181
warn!("NoisyDecoyProvider: failed to send decoy packet: {err:?}");
82-
} else if let Some(bytes) = body_bytes {
83-
try_replicate(&state, &manager, false, bytes).await;
82+
} else if let Some(body) = body_buf {
83+
try_replicate(&state, &manager, false, body).await;
8484
}
8585
}
8686
}

typhoon/src/flow/decoy/smooth.rs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -77,14 +77,14 @@ impl<T: IdentityType + Clone, AE: AsyncExecutor> SmoothDecoyProvider<T, AE> {
7777
if state_guard.try_spend_budget(decoy_length) {
7878
let decoy_packet = state_guard.create_decoy_packet(decoy_length, false);
7979
let should_rep = state_guard.should_replicate(false);
80+
let settings = Arc::clone(&state_guard.settings);
8081
drop(state_guard);
8182

82-
// Clone before consuming: Arc refcount bump (zero-copy) for optional replication.
83-
let body_bytes = should_rep.then(|| decoy_packet.slice_end(decoy_length).to_vec());
83+
let body_buf = should_rep.then(|| settings.pool().allocate_precise_from_slice_with_capacity(decoy_packet.slice_end(decoy_length), 0, 0));
8484
if let Err(err) = manager_arc.send_decoy_packet(decoy_packet).await {
8585
warn!("SmoothDecoyProvider: failed to send decoy packet: {err:?}");
86-
} else if let Some(bytes) = body_bytes {
87-
try_replicate(&state, &manager, false, bytes).await;
86+
} else if let Some(body) = body_buf {
87+
try_replicate(&state, &manager, false, body).await;
8888
}
8989
}
9090
}

typhoon/src/flow/decoy/sparse.rs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -77,14 +77,14 @@ impl<T: IdentityType + Clone, AE: AsyncExecutor> SparseDecoyProvider<T, AE> {
7777
if state_guard.try_spend_budget(decoy_length) {
7878
let decoy_packet = state_guard.create_decoy_packet(decoy_length, false);
7979
let should_rep = state_guard.should_replicate(false);
80+
let settings = Arc::clone(&state_guard.settings);
8081
drop(state_guard);
8182

82-
// Clone before consuming: Arc refcount bump (zero-copy) for optional replication.
83-
let body_bytes = should_rep.then(|| decoy_packet.slice_end(decoy_length).to_vec());
83+
let body_buf = should_rep.then(|| settings.pool().allocate_precise_from_slice_with_capacity(decoy_packet.slice_end(decoy_length), 0, 0));
8484
if let Err(err) = manager_arc.send_decoy_packet(decoy_packet).await {
8585
warn!("SparseDecoyProvider: failed to send decoy packet: {err:?}");
86-
} else if let Some(bytes) = body_bytes {
87-
try_replicate(&state, &manager, false, bytes).await;
86+
} else if let Some(body) = body_buf {
87+
try_replicate(&state, &manager, false, body).await;
8888
}
8989
}
9090
}

0 commit comments

Comments
 (0)